libeufin

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

commit 2b1907ef57b429d3e42fdeeba62298dfdb405643
parent 5376b61d95bc46857d001b47dd67bd410ed6f814
Author: Antoine A <>
Date:   Thu, 11 Jun 2026 14:01:20 +0200

ebisync: migrate the component

Diffstat:
MCargo.lock | 108++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
MCargo.toml | 20++++++++++++++------
MMakefile | 11++++++++++-
Adatabase-versioning/ebisync-0001.sql | 35+++++++++++++++++++++++++++++++++++
Adatabase-versioning/ebisync-drop.sql | 33+++++++++++++++++++++++++++++++++
Ddatabase-versioning/libeufin-ebisync-0001.sql | 35-----------------------------------
Ddatabase-versioning/libeufin-ebisync-drop.sql | 33---------------------------------
Mdebian/control | 13+++++++------
Mdebian/libeufin-ebisync.install | 4++--
Mdebian/rules | 9++++++++-
Mlibeufin-bank/Cargo.toml | 2+-
Mlibeufin-bank/src/config.rs | 2+-
Mlibeufin-bank/src/db.rs | 21+++++++++------------
Mlibeufin-bank/src/lib.rs | 14+++++++-------
Mlibeufin-ebics/Cargo.toml | 4+---
Mlibeufin-ebics/src/cli.rs | 10++++++++++
Mlibeufin-ebics/src/crypto.rs | 26+++++++++++++++-----------
Mlibeufin-ebics/src/db.rs | 64++++++++++++++++++++++++++++++++++++++++++++++++----------------
Mlibeufin-ebics/src/ebics.rs | 46++++++++++++++++++++++++++++++++++++----------
Mlibeufin-ebics/src/ebics/administrative.rs | 2+-
Mlibeufin-ebics/src/ebics/bts.rs | 6+++---
Mlibeufin-ebics/src/ebics/key_management.rs | 14+++++++++-----
Mlibeufin-ebics/src/keys.rs | 146++++++++++++++++++++++++++++++++++---------------------------------------------
Mlibeufin-ebics/src/setup.rs | 16++++++++--------
Mlibeufin-ebics/src/test.rs | 23++++++-----------------
Alibeufin-ebisync/Cargo.toml | 38++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/conf/test.conf | 15+++++++++++++++
Alibeufin-ebisync/ebisync.conf | 96+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/spa/index.html | 512+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/src/api.rs | 220+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/src/azure.rs | 219+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/src/config.rs | 250+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/src/db.rs | 60++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/src/lib.rs | 562+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Alibeufin-ebisync/src/main.rs | 27+++++++++++++++++++++++++++
Mlibeufin-nexus/Cargo.toml | 8+++-----
Mlibeufin-nexus/src/config.rs | 17++++++++++-------
Mlibeufin-nexus/src/db.rs | 59++++++++++-------------------------------------------------
Mlibeufin-nexus/src/fetch.rs | 9+++------
Mlibeufin-nexus/src/lib.rs | 57+++++++++++++++------------------------------------------
Mlibeufin-nexus/src/testing.rs | 10+++++-----
Mtestbench/Cargo.toml | 2+-
42 files changed, 2465 insertions(+), 393 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -217,6 +217,7 @@ dependencies = [ "matchit", "memchr", "mime", + "multer", "percent-encoding", "pin-project-lite", "serde_core", @@ -1312,6 +1313,31 @@ dependencies = [ ] [[package]] +name = "http-client" +version = "1.5.0" +source = "git+git://git.taler.net/taler-rust.git/#4e0704b3732a7a139f9ec7eef8a5adfdfa391aaf" +dependencies = [ + "compact_str", + "futures-util", + "http", + "http-body-util", + "hyper", + "hyper-rustls", + "hyper-util", + "rustls", + "serde", + "serde_json", + "serde_path_to_error", + "serde_urlencoded", + "taler-common", + "thiserror", + "tokio", + "tokio-util", + "tracing", + "url", +] + +[[package]] name = "http-range-header" version = "0.4.2" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1370,6 +1396,7 @@ dependencies = [ "hyper", "hyper-util", "rustls", + "rustls-platform-verifier", "tokio", "tokio-rustls", "tower-service", @@ -1786,6 +1813,39 @@ dependencies = [ ] [[package]] +name = "libeufin-ebisync" +version = "1.5.0" +dependencies = [ + "anyhow", + "aws-lc-rs", + "axum", + "clap", + "compact_str", + "dialoguer", + "http-client", + "jiff", + "libeufin-ebics", + "pretty_assertions", + "rand 0.10.1", + "reqwest", + "serde", + "serde_json", + "shlex", + "sqlx", + "taler-api", + "taler-build", + "taler-common", + "taler-macros", + "taler-test-utils", + "thiserror", + "tokio", + "tower-http", + "tracing", + "uuid", + "zip 8.6.0", +] + +[[package]] name = "libeufin-nexus" version = "1.5.0" dependencies = [ @@ -2005,6 +2065,23 @@ dependencies = [ ] [[package]] +name = "multer" +version = "3.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "83e87776546dc87511aa5ee218730c92b666d7264ab6ed41f9d215af9cd5224b" +dependencies = [ + "bytes", + "encoding_rs", + "futures-util", + "http", + "httparse", + "memchr", + "mime", + "spin", + "version_check", +] + +[[package]] name = "nix" version = "0.31.3" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2537,9 +2614,9 @@ dependencies = [ [[package]] name = "regex" -version = "1.12.3" +version = "1.12.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e10754a14b9137dd7b1e3e5b0493cc9171fdd105e0ab477f51b72e7f3ac0e276" +checksum = "f1292b7759ae1cb9ec195452d1390a074f0cd8541ab7a5a8c31cd6db45d4a6ba" dependencies = [ "aho-corasick", "memchr", @@ -2560,9 +2637,9 @@ dependencies = [ [[package]] name = "regex-syntax" -version = "0.8.10" +version = "0.8.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dc897dd8d9e8bd1ed8cdad82b5966c3e0ecae09fb1907d58efaa013543185d0a" +checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" [[package]] name = "reqwest" @@ -2707,6 +2784,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ef86cd5876211988985292b91c96a8f2d298df24e75989a43a3c73f2d4d8168b" dependencies = [ "aws-lc-rs", + "log", "once_cell", "rustls-pki-types", "rustls-webpki", @@ -3392,7 +3470,7 @@ dependencies = [ [[package]] name = "taler-api" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#25696769c34802040aa0af5a1022adfe0fd93f54" +source = "git+git://git.taler.net/taler-rust.git/#4e0704b3732a7a139f9ec7eef8a5adfdfa391aaf" dependencies = [ "aws-lc-rs", "axum", @@ -3419,12 +3497,12 @@ dependencies = [ [[package]] name = "taler-build" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#25696769c34802040aa0af5a1022adfe0fd93f54" +source = "git+git://git.taler.net/taler-rust.git/#4e0704b3732a7a139f9ec7eef8a5adfdfa391aaf" [[package]] name = "taler-common" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#25696769c34802040aa0af5a1022adfe0fd93f54" +source = "git+git://git.taler.net/taler-rust.git/#4e0704b3732a7a139f9ec7eef8a5adfdfa391aaf" dependencies = [ "anyhow", "aws-lc-rs", @@ -3453,7 +3531,7 @@ dependencies = [ [[package]] name = "taler-macros" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#25696769c34802040aa0af5a1022adfe0fd93f54" +source = "git+git://git.taler.net/taler-rust.git/#4e0704b3732a7a139f9ec7eef8a5adfdfa391aaf" dependencies = [ "proc-macro2", "quote", @@ -3463,7 +3541,7 @@ dependencies = [ [[package]] name = "taler-test-utils" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#25696769c34802040aa0af5a1022adfe0fd93f54" +source = "git+git://git.taler.net/taler-rust.git/#4e0704b3732a7a139f9ec7eef8a5adfdfa391aaf" dependencies = [ "aws-lc-rs", "axum", @@ -3940,9 +4018,9 @@ checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" [[package]] name = "uuid" -version = "1.23.2" +version = "1.23.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d258b83ceec21034727ecee8c382cfa6c3e133699b0742c64571814fb420c9f7" +checksum = "144d6b123cef80b301b8f72a9e2ca4370ddec21950d0a103dd22c437006d2db7" dependencies = [ "getrandom 0.4.2", "js-sys", @@ -4670,18 +4748,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.50" +version = "0.8.52" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3b065d4f0e55f82fae73202e189638116a87c55ab6b8e6c2721e13dd9d854ad1" +checksum = "ce1022995ff5ff5d841ad7d994facc23098cd40152f2c1d11cd607c6f530653f" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.50" +version = "0.8.52" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0b631b19d36a892ab55420c92dbc83ccd79274f25be714855d3074aa71cab639" +checksum = "1ae7f38b72ec2a254e2b87ef277cf2cd4fb97cbebf944faa6f33354da0867930" dependencies = [ "proc-macro2", "quote", diff --git a/Cargo.toml b/Cargo.toml @@ -1,6 +1,6 @@ [workspace] resolver = "3" -members = ["libeufin-bank", "libeufin-ebics", "libeufin-nexus", "testbench"] +members = ["libeufin-bank", "libeufin-ebics", "libeufin-ebisync", "libeufin-nexus", "testbench"] [workspace.package] version = "1.5.0" @@ -42,13 +42,20 @@ uuid = { version = "1.0", features = ["v4", "fast-rng", "serde"] } rand = "0.10" pretty_assertions = "1" dialoguer = "0.12" -#taler-common = { path = "../taler-rust/common/taler-common" } -#taler-api = { path = "../taler-rust/common/taler-api" } -#taler-build = { path = "../taler-rust/common/taler-build" } -#taler-test-utils = { path = "../taler-rust/common/taler-test-utils" } -#taler-macros = { path = "../taler-rust/common/taler-macros" } +zip = { version = "8.5", default-features = false, features = [ + "deflate-flate2-zlib-rs", +] } +tower-http = { version = "0.6", features = ["fs"]} +shlex = "2.0" taler-common = { git = "git://git.taler.net/taler-rust.git/" } taler-api = { git = "git://git.taler.net/taler-rust.git/" } taler-build = { git = "git://git.taler.net/taler-rust.git/" } taler-test-utils = { git = "git://git.taler.net/taler-rust.git/" } taler-macros = { git = "git://git.taler.net/taler-rust.git/" } +http-client = { git = "git://git.taler.net/taler-rust.git/" } +#taler-common = { path = "../taler-rust/common/taler-common" } +#taler-api = { path = "../taler-rust/common/taler-api" } +#taler-build = { path = "../taler-rust/common/taler-build" } +#taler-test-utils = { path = "../taler-rust/common/taler-test-utils" } +#taler-macros = { path = "../taler-rust/common/taler-macros" } +#http-client = { path = "../taler-rust/common/http-client" } +\ No newline at end of file diff --git a/Makefile b/Makefile @@ -13,7 +13,7 @@ all: build .PHONY: build build: - cargo build --release --bin libeufin-bank --bin libeufin-nexus + cargo build --release --bin libeufin-bank --bin libeufin-nexus --bin libeufin-ebisync .PHONY: install-nobuild-files install-nobuild-files: @@ -24,11 +24,16 @@ install-nobuild-files: install -m 644 -D -t $(share_dir)/libeufin/sql database-versioning/libeufin-bank*.sql install -m 644 -D -t $(share_dir)/libeufin/sql database-versioning/libeufin-nexus*.sql install -m 644 -D -t $(share_dir)/libeufin/sql database-versioning/libeufin-conversion*.sql + install -m 644 -D -t $(share_dir)/libeufin-ebisync/config.d libeufin-ebisync/ebisync.conf + install -m 644 -D -t $(share_dir)/libeufin-ebisync/sql database-versioning/versioning.sql + install -m 644 -D -t $(share_dir)/libeufin-ebisync/sql database-versioning/ebisync*.sql install -D -t $(bin_dir) contrib/libeufin-dbconfig install -D -t $(bin_dir) contrib/libeufin-ebisync-dbconfig install -D -t $(bin_dir) contrib/libeufin-tan-*.sh install -d $(share_dir)/libeufin/spa cp contrib/wallet-core/bank/* $(share_dir)/libeufin/spa/ + install -d $(share_dir)/libeufin-ebisync/spa + cp libeufin-ebisync/spa/* $(share_dir)/libeufin-ebisync/spa/ .PHONY: install install: build install-nobuild-files @@ -40,6 +45,10 @@ install: build install-nobuild-files install -D -t $(bin_dir) target/release/libeufin-nexus install -m 644 -D -t $(man_dir)/man1 doc/prebuilt/man/libeufin-nexus.1 install -m 644 -D -t $(man_dir)/man5 doc/prebuilt/man/libeufin-nexus.conf.5 +# Install libeufin-ebisync + install -D -t $(bin_dir) target/release/libeufin-ebisync + install -m 644 -D -t $(man_dir)/man1 doc/prebuilt/man/libeufin-ebisync.1 + install -m 644 -D -t $(man_dir)/man5 doc/prebuilt/man/libeufin-ebisync.conf.5 .PHONY: check check: install-nobuild-files diff --git a/database-versioning/ebisync-0001.sql b/database-versioning/ebisync-0001.sql @@ -0,0 +1,35 @@ +-- +-- This file is part of TALER +-- Copyright (C) 2025 Taler Systems SA +-- +-- TALER is free software; you can redistribute it and/or modify it under the +-- terms of the GNU General Public License as published by the Free Software +-- Foundation; either version 3, or (at your option) any later version. +-- +-- TALER is distributed in the hope that it will be useful, but WITHOUT ANY +-- WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR +-- A PARTICULAR PURPOSE. See the GNU General Public License for more details. +-- +-- You should have received a copy of the GNU General Public License along with +-- TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> + +BEGIN; + +SELECT _v.register_patch('ebisync-0001', NULL, NULL); + +CREATE SCHEMA ebisync; +SET search_path TO ebisync; + +CREATE TABLE kv ( + key TEXT NOT NULL PRIMARY KEY, + value JSONB NOT NULL +); +COMMENT ON TYPE kv + IS 'Store key/value data that do not fit well in a traditional relational table.'; + +CREATE TABLE pending_ebics_transactions ( + tx_id TEXT NOT NULL UNIQUE PRIMARY KEY +); +COMMENT ON TYPE pending_ebics_transactions + IS 'Store pending EBICS transactions ids to cleanly close them on failure.'; +COMMIT; diff --git a/database-versioning/ebisync-drop.sql b/database-versioning/ebisync-drop.sql @@ -0,0 +1,33 @@ +-- +-- This file is part of TALER +-- Copyright (C) 2025 Taler Systems SA +-- +-- TALER is free software; you can redistribute it and/or modify it under the +-- terms of the GNU General Public License as published by the Free Software +-- Foundation; either version 3, or (at your option) any later version. +-- +-- TALER is distributed in the hope that it will be useful, but WITHOUT ANY +-- WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR +-- A PARTICULAR PURPOSE. See the GNU General Public License for more details. +-- +-- You should have received a copy of the GNU General Public License along with +-- TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> + +BEGIN; + +DO +$do$ +DECLARE + patch text; +BEGIN + IF EXISTS(SELECT FROM information_schema.schemata WHERE schema_name='_v') THEN + FOR patch IN SELECT patch_name FROM _v.patches WHERE patch_name LIKE 'ebisync_%' LOOP + PERFORM _v.unregister_patch(patch); + END LOOP; + END IF; +END +$do$; + +DROP SCHEMA IF EXISTS ebisync CASCADE; + +COMMIT; diff --git a/database-versioning/libeufin-ebisync-0001.sql b/database-versioning/libeufin-ebisync-0001.sql @@ -1,35 +0,0 @@ --- --- This file is part of TALER --- Copyright (C) 2025 Taler Systems SA --- --- TALER is free software; you can redistribute it and/or modify it under the --- terms of the GNU General Public License as published by the Free Software --- Foundation; either version 3, or (at your option) any later version. --- --- TALER is distributed in the hope that it will be useful, but WITHOUT ANY --- WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR --- A PARTICULAR PURPOSE. See the GNU General Public License for more details. --- --- You should have received a copy of the GNU General Public License along with --- TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> - -BEGIN; - -SELECT _v.register_patch('libeufin-ebisync-0001', NULL, NULL); - -CREATE SCHEMA libeufin_ebisync; -SET search_path TO libeufin_ebisync; - -CREATE TABLE kv ( - key TEXT NOT NULL PRIMARY KEY, - value JSONB NOT NULL -); -COMMENT ON TYPE kv - IS 'Store key/value data that do not fit well in a traditional relational table.'; - -CREATE TABLE pending_ebics_transactions ( - tx_id TEXT NOT NULL UNIQUE PRIMARY KEY -); -COMMENT ON TYPE pending_ebics_transactions - IS 'Store pending EBICS transactions ids to cleanly close them on failure.'; -COMMIT; diff --git a/database-versioning/libeufin-ebisync-drop.sql b/database-versioning/libeufin-ebisync-drop.sql @@ -1,33 +0,0 @@ --- --- This file is part of TALER --- Copyright (C) 2025 Taler Systems SA --- --- TALER is free software; you can redistribute it and/or modify it under the --- terms of the GNU General Public License as published by the Free Software --- Foundation; either version 3, or (at your option) any later version. --- --- TALER is distributed in the hope that it will be useful, but WITHOUT ANY --- WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR --- A PARTICULAR PURPOSE. See the GNU General Public License for more details. --- --- You should have received a copy of the GNU General Public License along with --- TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> - -BEGIN; - -DO -$do$ -DECLARE - patch text; -BEGIN - IF EXISTS(SELECT FROM information_schema.schemata WHERE schema_name='_v') THEN - FOR patch IN SELECT patch_name FROM _v.patches WHERE patch_name LIKE 'libeufin_ebisync_%' LOOP - PERFORM _v.unregister_patch(patch); - END LOOP; - END IF; -END -$do$; - -DROP SCHEMA IF EXISTS libeufin_ebisync CASCADE; - -COMMIT; diff --git a/debian/control b/debian/control @@ -35,9 +35,10 @@ Recommends: postgresql (>= 14.0) Description: Software package to access a bank accounts via the EBICS protocol. -#Package: libeufin-ebisync -#Architecture: any -#Depends: ${misc:Depends}, ${shlibs:Depends} -#Recommends: -# postgresql (>= 14.0) -#Description: Software package to sync ISO20022 files via the EBICS protocol. +Package: libeufin-ebisync +Architecture: any +Depends: ${misc:Depends}, ${shlibs:Depends} +Recommends: + nginx | apache2 | httpd, + postgresql (>= 14.0) +Description: Software package to sync ISO20022 files via the EBICS protocol. diff --git a/debian/libeufin-ebisync.install b/debian/libeufin-ebisync.install @@ -6,9 +6,9 @@ target/release/libeufin-ebisync usr/bin/ contrib/libeufin-ebisync-dbconfig usr/bin/ database-versioning/versioning.sql usr/share/libeufin-ebisync/sql/ -database-versioning/libeufin-ebisync*.sql usr/share/libeufin-ebisync/sql/ +database-versioning/ebisync*.sql usr/share/libeufin-ebisync/sql/ -libeufin-ebisync/src/spa/* usr/share/libeufin-ebisync/spa +libeufin-ebisync/spa/* usr/share/libeufin-ebisync/spa libeufin-ebisync/ebisync.conf usr/share/libeufin-ebisync/config.d/ diff --git a/debian/rules b/debian/rules @@ -7,7 +7,7 @@ override_dh_auto_configure: rustup default stable override_dh_auto_build: - cargo build --release --bin libeufin-bank --bin libeufin-nexus + make build override_dh_auto_test: true @@ -15,10 +15,17 @@ override_dh_auto_test: override_dh_auto_install: true +# Override this step because it's very slow and likely +# unnecessary for us. +override_dh_strip_nondeterminism: + true + override_dh_installsystemd: dh_installsystemd --no-start --no-enable --no-stop-on-upgrade -p libeufin-bank --name libeufin-bank-gc dh_installsystemd --no-start --no-enable --no-stop-on-upgrade -p libeufin-bank --name libeufin-bank dh_installsystemd --no-start --no-enable --no-stop-on-upgrade -p libeufin-nexus --name libeufin-nexus-ebics-submit dh_installsystemd --no-start --no-enable --no-stop-on-upgrade -p libeufin-nexus --name libeufin-nexus-ebics-fetch dh_installsystemd --no-start --no-enable --no-stop-on-upgrade -p libeufin-nexus --name libeufin-nexus-httpd + dh_installsystemd --no-start --no-enable --no-stop-on-upgrade -p libeufin-ebisync --name libeufin-ebisync-fetch + dh_installsystemd --no-start --no-enable --no-stop-on-upgrade -p libeufin-ebisync --name libeufin-ebisync-httpd dh_installsystemd diff --git a/libeufin-bank/Cargo.toml b/libeufin-bank/Cargo.toml @@ -27,7 +27,7 @@ uuid.workspace = true axum.workspace = true rand.workspace = true dialoguer.workspace = true -tower-http = { version = "0.6", features = ["fs"]} +tower-http.workspace = true futures = "0.3" url = "2.5" regex = "1.12" diff --git a/libeufin-bank/src/config.rs b/libeufin-bank/src/config.rs @@ -192,7 +192,7 @@ impl BankCfg { gc_delete_after: s.span("gc_delete_after").require()?, pwd_check_quality: s.boolean("pwd_check").require()?, basic_auth_compat: s.boolean("pwd_auth_compat").require()?, - serve: Serve::parse(s)?, + serve: Serve::parse(&s)?, db_cfg: DbCfg::parse(cfg.section("libeufin-bankdb-postgres"))?, cfg, }) diff --git a/libeufin-bank/src/db.rs b/libeufin-bank/src/db.rs @@ -17,19 +17,18 @@ * <http://www.gnu.org/licenses/> */ +#![allow(clippy::too_many_arguments)] + use std::time::Duration; use jiff::Timestamp; use sqlx::{PgPool, postgres::PgRow}; -use taler_api::{db::TypeHelper, notification::NotificationChannel, serialized}; +use taler_api::{config::DbCfg, db::TypeHelper, notification::NotificationChannel, serialized}; use taler_common::types::amount::{Amount, Currency}; use tokio::join; use uuid::Uuid; -use crate::{ - api::{MonitorParams, MonitorResponse, withdrawal::WithdrawalStatus}, - config::BankCfg, -}; +use crate::api::{MonitorParams, MonitorResponse, withdrawal::WithdrawalStatus}; pub mod account; pub mod cashout; @@ -44,17 +43,15 @@ pub mod withdrawal; const SCHEMA: &str = "libeufin_bank"; -pub async fn pool(cfg: &BankCfg) -> anyhow::Result<PgPool> { - let db_cfg = &cfg.db_cfg; - let pool = taler_common::db::pool(db_cfg.cfg.clone(), SCHEMA).await?; +pub async fn pool(cfg: &DbCfg) -> anyhow::Result<PgPool> { + let pool = taler_common::db::pool(cfg.cfg.clone(), SCHEMA).await?; Ok(pool) } -pub async fn dbinit(cfg: &BankCfg, reset: bool) -> anyhow::Result<PgPool> { - let db_cfg = &cfg.db_cfg; - let pool = taler_common::db::pool(db_cfg.cfg.clone(), SCHEMA).await?; +pub async fn dbinit(cfg: &DbCfg, reset: bool) -> anyhow::Result<PgPool> { + let pool = taler_common::db::pool(cfg.cfg.clone(), SCHEMA).await?; let mut db = pool.acquire().await?; - taler_common::db::dbinit(&mut db, db_cfg.sql_dir.as_ref(), "libeufin-bank", reset).await?; + taler_common::db::dbinit(&mut db, cfg.sql_dir.as_ref(), "libeufin-bank", reset).await?; Ok(pool) } diff --git a/libeufin-bank/src/lib.rs b/libeufin-bank/src/lib.rs @@ -255,7 +255,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { match cmd { Cmd::Dbinit { reset } => { let cfg = BankCfg::parse(cfg)?; - let db = dbinit(&cfg, reset).await?; + let db = dbinit(&cfg.db_cfg, reset).await?; match create_admin_account(&db, &cfg, None).await? { CreationResult::Success(_) => { info!("Admin's account created") @@ -268,7 +268,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { } Cmd::Passwd { username, password } => { let cfg = BankCfg::parse(cfg)?; - let db = pool(&cfg).await?; + let db = pool(&cfg.db_cfg).await?; let pw = match password { Some(p) => p, @@ -299,7 +299,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { current_token, } => { let cfg = BankCfg::parse(cfg)?; - let db = pool(&cfg).await?; + let db = pool(&cfg.db_cfg).await?; let now = Timestamp::now(); let new = if let Some(current) = current_token && let Some(token) = access(&db, current.as_ref(), &now).await? @@ -338,7 +338,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { } Cmd::Serve => { let cfg = BankCfg::parse(cfg)?; - let db = pool(&cfg).await?; + let db = pool(&cfg.db_cfg).await?; if cfg.fiat.is_some() { info!("Ensure exchange account exists"); let Some(info) = bank_info(&db, &cfg.ctx, "exchange").await? else { @@ -405,7 +405,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { tan_channel, } => { let cfg = BankCfg::parse(cfg)?; - let db = pool(&cfg).await?; + let db = pool(&cfg.db_cfg).await?; let req = if let Some(json_str) = json { serde_json::from_str::<RegisterAccountRequest>(&json_str) .map_err(|e| anyhow!("Failed to parse JSON: {e}"))? @@ -466,7 +466,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { tan_channel, } => { let cfg = BankCfg::parse(cfg)?; - let db = pool(&cfg).await?; + let db = pool(&cfg.db_cfg).await?; let req = AccountReconfiguration { name, is_taler_exchange: exchange, @@ -508,7 +508,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { } Cmd::Gc => { let cfg = BankCfg::parse(cfg)?; - let db = pool(&cfg).await?; + let db = pool(&cfg.db_cfg).await?; let now = Zoned::now(); collect( &db, diff --git a/libeufin-ebics/Cargo.toml b/libeufin-ebics/Cargo.toml @@ -31,9 +31,7 @@ pretty_assertions.workspace = true dialoguer.workspace = true tempfile = "3" flate2 = { version = "1.0", features = ["zlib-rs"], default-features = false } -zip = { version = "8.5", default-features = false, features = [ - "deflate-flate2-zlib-rs", -] } +zip.workspace = true futures-util = "0.3" reqwest-websocket = "0.6.0" roxmltree = "0.21.1" diff --git a/libeufin-ebics/src/cli.rs b/libeufin-ebics/src/cli.rs @@ -25,3 +25,13 @@ pub struct EbicsLogs { #[arg(long("debug-ebics"), value_name("log_dir"), global(true))] pub dir: Option<PathBuf>, } + +#[derive(clap::Parser, Debug, Clone)] +pub struct EbicsArgs { + #[command(flatten)] + pub logs: EbicsLogs, + + /// Execute once and return, ignoring the 'FREQUENCY' configuration value + #[arg(long)] + pub transient: bool, +} diff --git a/libeufin-ebics/src/crypto.rs b/libeufin-ebics/src/crypto.rs @@ -38,8 +38,6 @@ use rcgen::{BasicConstraints, CertificateParams, DnType, IsCa, KeyUsagePurpose}; use taler_common::encoding::{base64, hex}; use x509_parser::prelude::{FromDer as _, X509Certificate}; -use crate::keys::RsaPub; - /// Generate a self-signed X.509 certificate from an RSA private key (PEM or DER) pub fn x509_certificate_from_rsa_private( pem: &str, @@ -79,22 +77,22 @@ pub fn x509_certificate_from_rsa_private( } /** Create an RSA public key from its components: [modulus] and [exponent] */ -pub fn rsa_pub_from_component(modulus: &[u8], exponent: &[u8]) -> anyhow::Result<RsaPub> { +pub fn rsa_pub_from_component(modulus: &[u8], exponent: &[u8]) -> anyhow::Result<PublicKey> { let key: PublicEncryptingKey = aws_lc_rs::rsa::PublicKeyComponents { n: modulus, e: exponent, } .try_into()?; - Ok(RsaPub::from_der(key.as_der()?.as_ref())?) + Ok(PublicKey::from_der(key.as_der()?.as_ref())?) } /// Extract an RSA public key from a X.509 certificate -pub fn rsa_private_from_b64_x509_certificate(encoded: &str) -> anyhow::Result<RsaPub> { +pub fn rsa_private_from_b64_x509_certificate(encoded: &str) -> anyhow::Result<PublicKey> { let der = base64::decode(encoded)?; let (_, cert) = X509Certificate::from_der(&der)?; let issuer_public_key = cert.public_key(); cert.verify_signature(Some(issuer_public_key))?; - Ok(RsaPub::from_der(issuer_public_key.raw)?) + Ok(PublicKey::from_der(issuer_public_key.raw)?) } /// Hash an RSA public key according to the EBICS standard (EBICS 2.5: 4.4.1.2.3). @@ -121,11 +119,14 @@ pub fn ebics_pub_key_hash(public_key: &PublicKey) -> Digest { ctx.finish() } -pub fn gen_ebics_e002_key(pub_key: PublicEncryptingKey) -> ([u8; 16], Vec<u8>) { +pub fn gen_ebics_e002_key(pub_key: &PublicKey) -> ([u8; 16], Vec<u8>) { let mut transaction_key = [0u8; 16]; SysRng.try_fill_bytes(&mut transaction_key).unwrap(); - let key = Pkcs1PublicEncryptingKey::new(pub_key).unwrap(); + let key = Pkcs1PublicEncryptingKey::new( + PublicEncryptingKey::from_der(pub_key.as_der().unwrap().as_ref()).unwrap(), + ) + .unwrap(); let mut encrypted_key = vec![0; key.ciphertext_size()]; key.encrypt(&transaction_key, &mut encrypted_key).unwrap(); @@ -211,7 +212,8 @@ pub fn verify_ebics_a006(sig: &[u8], data: &[u8], public_key_der: &PublicKey) -> #[cfg(test)] mod test { use aws_lc_rs::{ - rsa::{KeyPair, KeySize, PrivateDecryptingKey}, + encoding::AsDer, + rsa::{KeyPair, KeySize, PrivateDecryptingKey, PublicKey}, signature::KeyPair as _, }; use taler_common::encoding::hex; @@ -226,7 +228,9 @@ mod test { let data = b"Hello, World!"; let key = PrivateDecryptingKey::generate(KeySize::Rsa2048).unwrap(); - let (tx_key, encrypted_key) = gen_ebics_e002_key(key.public_key()); + let (tx_key, encrypted_key) = gen_ebics_e002_key( + &PublicKey::from_der(key.public_key().as_der().unwrap().as_ref()).unwrap(), + ); let enc = encrypt_ebics_e002(&tx_key, data.to_vec()); let key = decrypt_ebics_e002_key(key, &encrypted_key); let dec = decrypt_ebics_e002(&key, enc); @@ -273,7 +277,7 @@ mod test { &hex::decode(exponent).unwrap(), ) .unwrap(); - let hash = ebics_pub_key_hash(&key.key); + let hash = ebics_pub_key_hash(&key); assert_eq!(&hex::decode(expected).unwrap(), hash.as_ref()); } } diff --git a/libeufin-ebics/src/db.rs b/libeufin-ebics/src/db.rs @@ -15,37 +15,69 @@ */ use compact_str::CompactString; -use sqlx::{PgPool, Row, postgres::PgRow}; +use jiff::Timestamp; +use sqlx::{PgPool, Row, postgres::PgRow, types::Json}; +use taler_api::{db::BindHelper as _, serialized}; + +use crate::ebics::TaskStatus; /** Register a pending transaction */ pub async fn ebics_register(db: &PgPool, id: &str) -> sqlx::Result<()> { - sqlx::query( - "INSERT INTO pending_ebics_transactions (tx_id) VALUES ($1) ON CONFLICT DO NOTHING", - ) - .bind(id) - .execute(db) - .await?; + serialized!( + sqlx::query( + "INSERT INTO pending_ebics_transactions (tx_id) VALUES ($1) ON CONFLICT DO NOTHING", + ) + .bind(id) + .execute(db) + )?; Ok(()) } /** Register a pending transaction */ pub async fn ebics_remove(db: &PgPool, id: &str) -> sqlx::Result<()> { - sqlx::query("DELETE FROM pending_ebics_transactions WHERE tx_id = $1") - .bind(id) - .execute(db) - .await?; + serialized!( + sqlx::query("DELETE FROM pending_ebics_transactions WHERE tx_id = $1") + .bind(id) + .execute(db) + )?; Ok(()) } /** Register a pending transaction */ pub async fn ebics_first(db: &PgPool) -> sqlx::Result<Option<CompactString>> { - sqlx::query("SELECT tx_id FROM pending_ebics_transactions LIMIT 1") - .try_map(|r: PgRow| r.try_get(0)) - .fetch_optional(db) - .await + serialized!( + sqlx::query("SELECT tx_id FROM pending_ebics_transactions LIMIT 1") + .try_map(|r: PgRow| r.try_get(0)) + .fetch_optional(db) + ) +} + +/** Get current value for [key] */ +pub async fn get_task_status(db: &PgPool, key: &str) -> sqlx::Result<Option<TaskStatus>> { + serialized!( + sqlx::query_scalar::<_, Json<TaskStatus>>("SELECT value FROM kv WHERE key=$1") + .bind(key) + .fetch_optional(db) + ) + .map(|it| it.map(|it| it.0)) +} + +/** Update a TaskStatus timestamp */ +pub async fn update_task_status( + db: &PgPool, + key: &str, + timestamp: &Timestamp, + success: bool, +) -> sqlx::Result<()> { + serialized!( + sqlx::query(if success { + "INSERT INTO kv (key, value) VALUES ($1, jsonb_build_object('last_successfull', $2, 'last_trial', $3)) ON CONFLICT (key) DO UPDATE SET value=EXCLUDED.value" + } else { + "INSERT INTO kv (key, value) VALUES ($1, jsonb_build_object('last_trial', $2)) ON CONFLICT (key) DO UPDATE SET value=jsonb_set(EXCLUDED.value, '{last_trial}'::text[], to_jsonb($3))" + }).bind(key).bind_timestamp(timestamp).bind_timestamp(timestamp).execute(db))?; + Ok(()) } -#[cfg(test)] pub mod test { use sqlx::PgPool; diff --git a/libeufin-ebics/src/ebics.rs b/libeufin-ebics/src/ebics.rs @@ -19,7 +19,11 @@ use std::{borrow::Cow, io::Write as _}; -use aws_lc_rs::{digest::Digest, rsa::PrivateDecryptingKey}; +use aws_lc_rs::{ + digest::Digest, + encoding::AsDer, + rsa::{self, PrivateDecryptingKey}, +}; use compact_str::CompactString; use flate2::write::ZlibDecoder; use jiff::Timestamp; @@ -28,6 +32,7 @@ use reqwest::{ Client, ClientBuilder, StatusCode, header::{CONTENT_TYPE, HeaderValue}, }; +use serde::{Deserialize, Deserializer, Serialize, Serializer}; use sqlx::PgPool; use taler_common::encoding::base64; use tracing::{debug, info, trace, warn}; @@ -495,12 +500,12 @@ impl<'a> EbicsClient<'a> { client: &ClientKeys, bank: &BankKeys, order: &Order, - payload: &str, + payload: &[u8], ) -> Result<CompactString, EbicsError> { debug!(target: "ebics", "Uploading order {order}"); let mut ctx = EbicsCtx::new(order); - self.logger.log_payload(&ctx, payload.as_bytes(), "xml")?; + self.logger.log_payload(&ctx, payload, "xml")?; let payload = prepare_upload_payload(&self.cfg, client, bank, payload); // Init phase @@ -542,12 +547,14 @@ impl PreparedUploadData { /** Decrypts and decompresses EBICS BTS payload */ fn decrypt_and_decompress_payload( - client_encryption_key: &PrivateDecryptingKey, + key_pair: &rsa::KeyPair, encryption_info: DataEncryptionInfo, segments: Vec<Vec<u8>>, ) -> Vec<u8> { + let client_encryption_key = + PrivateDecryptingKey::from_pkcs8(key_pair.as_der().unwrap().as_ref()).unwrap(); // TODO check bank_pub_digest - let tx_key = decrypt_ebics_e002_key(client_encryption_key.clone(), &encryption_info.tx_key); + let tx_key = decrypt_ebics_e002_key(client_encryption_key, &encryption_info.tx_key); let mut decoder = ZlibDecoder::new(Vec::new()); for segment in segments { let decrypted = decrypt_ebics_e002(&tx_key, segment); @@ -561,12 +568,12 @@ fn prepare_upload_payload( cfg: &EbicsHostCfg, client: &ClientKeys, bank: &BankKeys, - payload: &str, + payload: &[u8], ) -> PreparedUploadData { - let digest = digest_ebics_order_a006(payload.as_bytes()); + let digest = digest_ebics_order_a006(payload); // Generate ephemeral transaction key - let (tx_key, encrypted_key) = gen_ebics_e002_key(bank.enc.enc.clone()); + let (tx_key, encrypted_key) = gen_ebics_e002_key(&bank.enc); // Compress and encrypt order signature let signature_data = { @@ -588,7 +595,7 @@ fn prepare_upload_payload( // Compress and encrypt payload let payload = { - let deflated = deflate(payload.as_bytes()); + let deflated = deflate(payload); let encrypted = encrypt_ebics_e002(&tx_key, deflated); base64::encode(encrypted) }; @@ -679,7 +686,7 @@ pub async fn tx_check( .take(2000000) .map(char::from) .collect(); - let payload = prepare_upload_payload(&ebics.cfg, client, bank, &random_string); + let payload = prepare_upload_payload(&ebics.cfg, client, bank, random_string.as_bytes()); match ebics .post_bts( u_init(&ebics.cfg, bank, client, submit, &payload), @@ -755,3 +762,22 @@ pub async fn tx_check( Ok(result) } + +#[derive(Debug, Serialize, Deserialize, Clone, Default)] +pub struct TaskStatus { + #[serde(serialize_with = "ser_micros", deserialize_with = "de_micros", default)] + pub last_successfull: Option<Timestamp>, + #[serde(serialize_with = "ser_micros", deserialize_with = "de_micros", default)] + pub last_trial: Option<Timestamp>, +} + +fn ser_micros<S: Serializer>(key: &Option<Timestamp>, serializer: S) -> Result<S::Ok, S::Error> { + key.map(|it| it.as_microsecond()).serialize(serializer) +} + +fn de_micros<'de, D: Deserializer<'de>>(deserializer: D) -> Result<Option<Timestamp>, D::Error> { + Option::<i64>::deserialize(deserializer)? + .map(Timestamp::from_microsecond) + .transpose() + .map_err(|e| serde::de::Error::custom(e.to_string())) +} diff --git a/libeufin-ebics/src/ebics/administrative.rs b/libeufin-ebics/src/ebics/administrative.rs @@ -57,7 +57,7 @@ pub struct HKD { pub struct PartnerInfo { pub name: Option<CompactString>, pub accounts: Box<[AccountInfo]>, - pub orders: Box<[OrderInfo]>, + pub orders: Vec<OrderInfo>, } pub struct OrderInfo { pub order: Order, diff --git a/libeufin-ebics/src/ebics/bts.rs b/libeufin-ebics/src/ebics/bts.rs @@ -59,8 +59,8 @@ fn signed_request( fn bank_digest(w: &mut XmlWriter, bank: &BankKeys) { xml!(w => "BankPubKeyDigests" { - "Authentication" "Version"="X002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256" : base64::fmt(ebics_pub_key_hash(&bank.auth.key)), - "Encryption" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256" : base64::fmt(ebics_pub_key_hash(&bank.enc.key)) + "Authentication" "Version"="X002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256" : base64::fmt(ebics_pub_key_hash(&bank.auth)), + "Encryption" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256" : base64::fmt(ebics_pub_key_hash(&bank.enc)) }, "SecurityMedium": "0000" ) @@ -241,7 +241,7 @@ pub fn u_init( "body" { "DataTransfer" { "DataEncryptionInfo" "authenticate"="true" { - "EncryptionPubKeyDigest" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256": base64::fmt(ebics_pub_key_hash(&bank.enc.key)), + "EncryptionPubKeyDigest" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256": base64::fmt(ebics_pub_key_hash(&bank.enc)), "TransactionKey": base64::fmt(&data.encrypted_key) }, "SignatureData" "authenticate"="true" : data.signature_data, diff --git a/libeufin-ebics/src/ebics/key_management.rs b/libeufin-ebics/src/ebics/key_management.rs @@ -20,7 +20,10 @@ use std::{borrow::Cow, io::Write as _}; use anyhow::bail; -use aws_lc_rs::encoding::{AsDer, Pkcs8V1Der}; +use aws_lc_rs::{ + encoding::{AsDer, Pkcs8V1Der}, + rsa::PublicKey, +}; use compact_str::CompactStringExt; use flate2::{Compression, write::ZlibEncoder}; use taler_common::encoding::base64; @@ -34,7 +37,7 @@ use crate::{ bts::DataEncryptionInfo, decrypt_and_decompress_payload, ebics_code::EbicsReturnCode, order::Order, }, - keys::{self, BankKeys, ClientKeys, RsaPub}, + keys::{self, BankKeys, ClientKeys}, xml, xml::{Xml, XmlAccess as _, XmlWriter}, xml_sign::sign_ebics, @@ -68,8 +71,9 @@ impl EbicsClient<'_> { Order::HIA => client.submitted_hia = true, _ => unreachable!("Only INI & HIA are supported for client keys"), } - keys::persist_client_keys(client, cfg.client.as_ref()).ctx(&ctx)?; - // TODO better error: Could not update the $order state on disk + keys::persist_client_keys(client, cfg.client.as_ref()) + .map_err(|e| EbicsErrKind::Custom(e.to_string().into())) + .ctx(&ctx)?; Ok(()) } @@ -250,7 +254,7 @@ impl EbicsClient<'_> { } } -pub fn rsa_pub_key(xml: Xml) -> xml::Result<RsaPub> { +pub fn rsa_pub_key(xml: Xml) -> xml::Result<PublicKey> { xml.one("X509Data") .one("X509Certificate") .decode(rsa_private_from_b64_x509_certificate) diff --git a/libeufin-ebics/src/keys.rs b/libeufin-ebics/src/keys.rs @@ -19,13 +19,10 @@ use std::{borrow::Cow, io::ErrorKind, path::Path}; -use anyhow::bail; +use anyhow::{anyhow, bail}; use aws_lc_rs::{ - encoding::{AsDer, Pkcs8V1Der}, - error::KeyRejected, - rsa::{ - KeyPair, KeySize, PrivateDecryptingKey, PublicEncryptingKey, PublicKey, PublicKeyComponents, - }, + encoding::{AsDer, Pkcs8V1Der, PublicKeyX509Der}, + rsa::{KeyPair, KeySize, PublicKey}, signature::KeyPair as _, }; use serde::{Deserialize, Deserializer, Serialize, Serializer}; @@ -47,9 +44,9 @@ pub struct ClientKeys { #[serde( rename = "encryption_private_key", serialize_with = "ser_pkcs8", - deserialize_with = "de_ras_priv_base32" + deserialize_with = "de_ras_sign_base32" )] - pub enc: PrivateDecryptingKey, + pub enc: KeyPair, #[serde( rename = "authentication_private_key", serialize_with = "ser_pkcs8", @@ -64,7 +61,7 @@ impl ClientKeys { pub fn generate() -> anyhow::Result<Self> { Ok(Self { sign: KeyPair::generate(KeySize::Rsa2048)?, - enc: PrivateDecryptingKey::generate(KeySize::Rsa2048)?, + enc: KeyPair::generate(KeySize::Rsa2048)?, auth: KeyPair::generate(KeySize::Rsa2048)?, submitted_ini: false, submitted_hia: false, @@ -72,80 +69,42 @@ impl ClientKeys { } } -#[derive(Debug)] -pub struct RsaPub { - pub enc: PublicEncryptingKey, - pub key: PublicKey, -} - -impl RsaPub { - pub fn generate() -> Self { - let key = KeyPair::generate(KeySize::Rsa2048).unwrap(); - Self::from_der(key.public_key().as_der().unwrap().as_ref()).unwrap() - } - - pub fn from_der(der: &[u8]) -> Result<Self, KeyRejected> { - let key = PublicKey::from_der(der)?; - let component = PublicKeyComponents { - n: key.modulus().big_endian_without_leading_zero(), - e: key.exponent().big_endian_without_leading_zero(), - }; - let enc = component.try_into().map_err(|_| KeyRejected::from(()))?; - Ok(Self { enc, key }) - } -} - -impl serde::Serialize for RsaPub { - fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> - where - S: Serializer, - { - let der = self - .key - .as_der() - .map_err(|e| serde::ser::Error::custom(e.to_string()))?; - let base32 = base32::encode(der.as_ref()); - base32.serialize(serializer) - } -} - -impl<'de> serde::Deserialize<'de> for RsaPub { - fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> - where - D: Deserializer<'de>, - { - let base32 = Cow::<str>::deserialize(deserializer)?; - let der = base32::decode(base32.as_bytes()) - .map_err(|e| serde::de::Error::custom(e.to_string()))?; - Self::from_der(&der).map_err(|e| serde::de::Error::custom(e.to_string())) - } +#[derive(Debug, Serialize, Deserialize)] +pub struct BankKeys { + #[serde( + rename = "bank_encryption_public_key", + serialize_with = "ser_der", + deserialize_with = "de_ras_priv_base32" + )] + pub enc: PublicKey, + #[serde( + rename = "bank_authentication_public_key", + serialize_with = "ser_der", + deserialize_with = "de_ras_priv_base32" + )] + pub auth: PublicKey, + pub accepted: bool, } -impl PartialEq for RsaPub { +impl PartialEq for BankKeys { fn eq(&self, other: &Self) -> bool { - self.key.exponent().big_endian_without_leading_zero() - == other.key.exponent().big_endian_without_leading_zero() - && self.key.modulus().big_endian_without_leading_zero() - == other.key.modulus().big_endian_without_leading_zero() + self.enc.as_ref() == other.enc.as_ref() && self.auth.as_ref() == other.auth.as_ref() } } -impl Eq for RsaPub {} - -#[derive(Debug, Serialize, Deserialize, PartialEq, Eq)] -pub struct BankKeys { - #[serde(rename = "bank_encryption_public_key")] - pub enc: RsaPub, - #[serde(rename = "bank_authentication_public_key")] - pub auth: RsaPub, - pub accepted: bool, -} +impl Eq for BankKeys {} impl BankKeys { pub fn generate() -> Self { Self { - enc: RsaPub::generate(), - auth: RsaPub::generate(), + enc: KeyPair::generate(KeySize::Rsa2048) + .unwrap() + .public_key() + .clone(), + auth: KeyPair::generate(KeySize::Rsa2048) + .unwrap() + .public_key() + .clone(), accepted: false, } } @@ -163,15 +122,26 @@ where base32.serialize(serializer) } -fn de_ras_priv_base32<'de, D>(deserializer: D) -> Result<PrivateDecryptingKey, D::Error> +fn ser_der<S, K>(key: &K, serializer: S) -> Result<S::Ok, S::Error> +where + K: AsDer<PublicKeyX509Der<'static>>, + S: Serializer, +{ + let der = key + .as_der() + .map_err(|e| serde::ser::Error::custom(e.to_string()))?; + let base32 = base32::encode(der.as_ref()); + base32.serialize(serializer) +} + +fn de_ras_priv_base32<'de, D>(deserializer: D) -> Result<PublicKey, D::Error> where D: Deserializer<'de>, { let base32 = Cow::<str>::deserialize(deserializer)?; let der = base32::decode(base32.as_bytes()).map_err(|e| serde::de::Error::custom(e.to_string()))?; - let key = PrivateDecryptingKey::from_pkcs8(&der) - .map_err(|e| serde::de::Error::custom(e.to_string()))?; + let key = PublicKey::from_der(&der).map_err(|e| serde::de::Error::custom(e.to_string()))?; Ok(key) } @@ -187,16 +157,24 @@ where } /// Persist the bank keys file to disk -pub fn persist_bank_keys(keys: &BankKeys, location: &Path) -> std::io::Result<()> { - json_file::persist(location, keys)?; - // TODO better error message "bank public keys" - Ok(()) +pub fn persist_bank_keys(keys: &BankKeys, path: &Path) -> anyhow::Result<()> { + json_file::persist(path, keys).map_err(|e| { + anyhow!( + "Could not write bank public keys at '{}': {}", + path.to_string_lossy(), + e.kind() + ) + }) } -pub fn persist_client_keys(keys: &ClientKeys, location: &Path) -> std::io::Result<()> { - json_file::persist(location, keys)?; - // TODO better error message "client private keys" - Ok(()) +pub fn persist_client_keys(keys: &ClientKeys, path: &Path) -> anyhow::Result<()> { + json_file::persist(path, keys).map_err(|e| { + anyhow!( + "Could not write client private keys at '{}': {}", + path.to_string_lossy(), + e.kind() + ) + }) } /// Load the bank keys file from disk diff --git a/libeufin-ebics/src/setup.rs b/libeufin-ebics/src/setup.rs @@ -105,19 +105,19 @@ pub async fn ebics_setup( let new = ebics.hpb(&client).await?; if let Some(current) = bank { // Check current bank keys - if current.enc != new.enc { + if current.enc.as_ref() != new.enc.as_ref() { bail!( "On disk bank encryption key stored at {} doesn't match server key\nDisk: {}\nServer: {}", cfg.keys.bank, - hex_chunk_by_two(ebics_pub_key_hash(&current.enc.key)), - hex_chunk_by_two(ebics_pub_key_hash(&new.enc.key)) + hex_chunk_by_two(ebics_pub_key_hash(&current.enc)), + hex_chunk_by_two(ebics_pub_key_hash(&new.enc)) ) - } else if current.auth != new.auth { + } else if current.auth.as_ref() != new.auth.as_ref() { bail!( "On disk bank authentication key stored at {} doesn't match server key\nDisk: {}\nServer: {}", cfg.keys.bank, - hex_chunk_by_two(ebics_pub_key_hash(&current.auth.key)), - hex_chunk_by_two(ebics_pub_key_hash(&new.auth.key)) + hex_chunk_by_two(ebics_pub_key_hash(&current.auth)), + hex_chunk_by_two(ebics_pub_key_hash(&new.auth)) ) } } else { @@ -128,8 +128,8 @@ pub async fn ebics_setup( let mut bank = new; if !bank.accepted { // Finishing the setup by accepting the bank keys. - let enc_hash = ebics_pub_key_hash(&bank.enc.key); - let auth_hash = ebics_pub_key_hash(&bank.auth.key); + let enc_hash = ebics_pub_key_hash(&bank.enc); + let auth_hash = ebics_pub_key_hash(&bank.auth); if auto_accept_keys { bank.accepted = true } else if let Some(enc) = cfg.enc diff --git a/libeufin-ebics/src/test.rs b/libeufin-ebics/src/test.rs @@ -24,10 +24,7 @@ use std::{ time::Duration, }; -use aws_lc_rs::{ - encoding::AsDer, - rsa::{KeyPair, KeySize, PublicEncryptingKey, PublicKey}, -}; +use aws_lc_rs::rsa::{KeyPair, KeySize, PublicKey}; use axum::{body::Bytes, response::IntoResponse as _, routing::post}; use compact_str::CompactString; use jiff::{ @@ -141,11 +138,7 @@ impl EbicsState { fn ebics_response_payload(&mut self, payload: &str, last: bool) -> EbicsRes { let tx_id = self.tx_id.insert(rand_ebics_id()); let deflated = deflate(payload.as_bytes()); - let client_enc = PublicEncryptingKey::from_der( - self.client_enc.as_ref().unwrap().as_der().unwrap().as_ref(), - ) - .unwrap(); - let (tx_key, encrypted_key) = gen_ebics_e002_key(client_enc); + let (tx_key, encrypted_key) = gen_ebics_e002_key(self.client_enc.as_ref().unwrap()); let encrypted = encrypt_ebics_e002(&tx_key, deflated); let xml = xml!("ebicsResponse" "xmlns"="http://www.ebics.org/H005" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" { "header" "authenticate"="true" { @@ -214,7 +207,7 @@ impl EbicsState { Self::parse_unsecure_request(body, "INI", "SignaturePubKeyOrderData", |root| { let n = root.one("SignaturePubKeyInfo")?; assert_eq!(n.one("SignatureVersion")?.text(), "A006"); - self.client_sign = Some(rsa_pub_key(n)?.key); + self.client_sign = Some(rsa_pub_key(n)?); Ok(()) }); EbicsRes::Ok( @@ -236,11 +229,11 @@ impl EbicsState { Self::parse_unsecure_request(body, "HIA", "HIARequestOrderData", |root| { let n = root.one("AuthenticationPubKeyInfo")?; assert_eq!(n.one("AuthenticationVersion")?.text(), "X002"); - self.client_auth = Some(rsa_pub_key(n)?.key); + self.client_auth = Some(rsa_pub_key(n)?); let n = root.one("EncryptionPubKeyInfo")?; assert_eq!(n.one("EncryptionVersion")?.text(), "E002"); - self.client_enc = Some(rsa_pub_key(n)?.key); + self.client_enc = Some(rsa_pub_key(n)?); Ok(()) }); EbicsRes::Ok( @@ -283,11 +276,7 @@ impl EbicsState { } }); let deflated = deflate(payload.as_bytes()); - let client_enc = PublicEncryptingKey::from_der( - self.client_enc.as_ref().unwrap().as_der().unwrap().as_ref(), - ) - .unwrap(); - let (tx_key, encrypted_key) = gen_ebics_e002_key(client_enc); + let (tx_key, encrypted_key) = gen_ebics_e002_key(self.client_enc.as_ref().unwrap()); let encrypted = encrypt_ebics_e002(&tx_key, deflated); EbicsRes::Ok( xml!("ebicsKeyManagementResponse" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" "xmlns"="http://www.ebics.org/H005" { diff --git a/libeufin-ebisync/Cargo.toml b/libeufin-ebisync/Cargo.toml @@ -0,0 +1,37 @@ +[package] +name = "libeufin-ebisync" +version.workspace = true +edition.workspace = true +authors.workspace = true +homepage.workspace = true +repository.workspace = true +license-file.workspace = true + +[dependencies] +libeufin-ebics.workspace = true +tower-http.workspace = true +compact_str.workspace = true +aws-lc-rs.workspace = true +reqwest.workspace = true +sqlx.workspace = true +thiserror.workspace = true +tracing.workspace = true +serde.workspace = true +serde_json.workspace = true +taler-common.workspace = true +taler-api.workspace = true +taler-build.workspace = true +taler-macros.workspace = true +taler-test-utils.workspace = true +zip.workspace = true +axum = {workspace = true, features = ["multipart"]} +anyhow.workspace = true +tokio.workspace = true +jiff.workspace = true +clap.workspace = true +uuid.workspace = true +rand.workspace = true +pretty_assertions.workspace = true +dialoguer.workspace = true +http-client.workspace = true +shlex.workspace = true +\ No newline at end of file diff --git a/libeufin-ebisync/conf/test.conf b/libeufin-ebisync/conf/test.conf @@ -0,0 +1,14 @@ +[ebisync] +HOST_BASE_URL = http://localhost:8080/ebicsweb +BANK_PUBLIC_KEYS_FILE = /tmp/ebics-test/bank-keys.json +CLIENT_PRIVATE_KEYS_FILE = /tmp/ebics-test/client-keys.json +HOST_ID = PFEBICS +USER_ID = PFC00563 +PARTNER_ID = PFC00563 + +[ebisyncdb-postgres] +CONFIG = postgresql:///libeufincheck + +[ebisync-submit] +SOURCE = ebisync-api +AUTH_METHOD = none +\ No newline at end of file diff --git a/libeufin-ebisync/ebisync.conf b/libeufin-ebisync/ebisync.conf @@ -0,0 +1,95 @@ +[paths] +EBISYNC_HOME = /var/lib/libeufin-ebisync/ + +[ebisync] +# Base URL of the bank server. +HOST_BASE_URL = + +# EBICS host ID. +HOST_ID = + +# EBICS user ID, as assigned by the bank. +USER_ID = + +# EBICS partner ID, as assigned by the bank. +PARTNER_ID = + +# EBICS partner ID, as assigned by the bank. +SYSTEM_ID = + +# File that holds the bank EBICS keys. +BANK_PUBLIC_KEYS_FILE = ${EBISYNC_HOME}/bank-ebics-keys.json + +# File that holds the client EBICS keys. +CLIENT_PRIVATE_KEYS_FILE = ${EBISYNC_HOME}/client-ebics-keys.json + +[ebisync-setup] +# Bank encryption public key hash +# BANK_ENCRYPTION_PUB_KEY_HASH = + +# Bank authentication public key hash +# BANK_AUTHENTICATION_PUB_KEY_HASH = + +[ebisync-fetch] +# How often should ebics-fetch run when the bank does not support real time notification +FREQUENCY = 30m + +# At what time of day should ebics-fetch perform a checkpoint +CHECKPOINT_TIME_OF_DAY = 19:00 + +# Where should the ebics file be stored? This his can either can be azure-blob-storage or none +DESTINATION = none + +# Azure API account base url for azure-blob-storage +# AZURE_API_URL = https://myaccount.blob.core.windows.net/ + +# Azure API account name for azure-blob-storage +# AZURE_ACCOUNT_NAME = myaccount + +# Azure API account key for azure-blob-storage +# AZURE_ACCOUNT_KEY = Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw== + +# Which Azure Blob Storage container to use for azure-blob-storage +# AZURE_COUNTAINER = mycontainer + +[ebisync-submit] +# Where does the ebics file come from? This his can either can be ebisync-api or none +SOURCE = none + +# Authentication scheme used by the API, this can either can be basic, bearer or none. +# AUTH_METHOD = bearer + +# User name for basic authentication scheme +# USERNAME = + +# Password for basic authentication scheme +# PASSWORD = + +# Token for bearer authentication scheme +# TOKEN = + +[ebisync-httpd] +# How "libeufin-ebisync serve" serves its API, this can either be tcp or unix +SERVE = tcp + +# Port on which the HTTP server listens, e.g. 9967. Only used if SERVE is tcp. +PORT = 8080 + +# Which IP address should we bind to? E.g. ``127.0.0.1`` or ``::1``for loopback. Can also be given as a hostname. Only used if SERVE is tcp. +BIND_TO = 0.0.0.0 + +# Which unix domain path should we bind to? Only used if SERVE is unix. +# UNIXPATH = libeufin-ebisync.sock + +# What should be the file access permissions for UNIXPATH? Only used if SERVE is unix. +# UNIXPATH_MODE = 660 + +# Path to spa files +SPA = $DATADIR/spa + +[ebisyncdb-postgres] +# Where are the SQL files to setup our tables? +SQL_DIR = $DATADIR/sql/ + +# DB connection string +CONFIG = postgres:///libeufin-ebisync +\ No newline at end of file diff --git a/libeufin-ebisync/spa/index.html b/libeufin-ebisync/spa/index.html @@ -0,0 +1,511 @@ +<!DOCTYPE html> +<html lang="en"> +<head> + <meta charset="UTF-8"> + <meta name="viewport" content="width=device-width, initial-scale=1.0"> + <title>LibEuFin EbiSync - File Submission Portal</title> + <link rel="preconnect" href="https://fonts.googleapis.com"> + <link rel="preconnect" href="https://fonts.gstatic.com" crossorigin> + <link href="https://fonts.googleapis.com/css2?family=Cormorant+Garamond:wght@300;400;600;700&family=Montserrat:wght@300;400;500;600&display=swap" rel="stylesheet"> + <style> + :root { + --cream: #FAF8F3; + --charcoal: #2B2B2B; + --rust: #C1503D; + --sage: #8B9A7E; + --gold: #D4AF37; + --shadow: rgba(43, 43, 43, 0.08); + } + + * { + margin: 0; + padding: 0; + box-sizing: border-box; + } + + body { + font-family: 'Montserrat', sans-serif; + background: var(--cream); + color: var(--charcoal); + line-height: 1.7; + overflow-x: hidden; + } + + /* Animated background texture */ + body::before { + content: ''; + position: fixed; + top: 0; + left: 0; + width: 100%; + height: 100%; + background-image: + repeating-linear-gradient(90deg, transparent, transparent 2px, rgba(43, 43, 43, 0.015) 2px, rgba(43, 43, 43, 0.015) 4px), + repeating-linear-gradient(0deg, transparent, transparent 2px, rgba(43, 43, 43, 0.015) 2px, rgba(43, 43, 43, 0.015) 4px); + pointer-events: none; + z-index: 0; + } + + .container { + max-width: 1200px; + margin: 0 auto; + padding: 0 2rem; + position: relative; + z-index: 1; + } + + header { + padding: 2rem 0 2rem; + } + + h1 { + font-family: 'Cormorant Garamond', serif; + font-size: 4.5rem; + font-weight: 300; + letter-spacing: -0.02em; + line-height: 1.1; + margin-bottom: 1rem; + } + + .subtitle { + font-size: 1rem; + letter-spacing: 0.15em; + text-transform: uppercase; + color: var(--rust); + font-weight: 500; + } + + .version-text { + margin-left: auto; + font-size: 0.75rem; + color: rgba(43, 43, 43, 0.5); + font-family: 'Courier New', monospace; + } + + main { + display: grid; + grid-template-columns: 1fr 1fr; + gap: 3rem; + margin-bottom: 6rem; + } + + section { + background: white; + padding: 1rem 2rem; + border: 1px solid rgba(43, 43, 43, 0.1); + position: relative; + transition: all 0.4s ease; + } + + section::before { + content: ''; + position: absolute; + top: 0; + left: 0; + width: 4px; + height: 100%; + background: var(--rust); + transform: scaleY(0); + transition: transform 0.4s ease; + } + + section:hover::before { + transform: scaleY(1); + } + + h2 { + font-family: 'Cormorant Garamond', serif; + font-size: 2.5rem; + font-weight: 400; + margin-bottom: 0.5rem; + letter-spacing: -0.01em; + } + + .order-card { + padding: 1.5rem; + margin-bottom: 1rem; + border: 1px solid rgba(43, 43, 43, 0.1); + cursor: pointer; + transition: all 0.3s ease; + position: relative; + background: var(--cream); + } + + .order-card:hover { + transform: translateX(8px); + border-color: var(--rust); + box-shadow: -8px 0 0 var(--rust); + } + + .order-card.selected { + background: var(--charcoal); + color: var(--cream); + border-color: var(--charcoal); + } + + .order-card.selected .order-id { + color: var(--gold); + } + + .order-id { + font-family: 'Courier New', monospace; + color: var(--rust); + margin-bottom: 0.5rem; + font-weight: bold; + } + + .order-description { + font-size: 0.875rem; + line-height: 1.6; + } + + .file-upload-zone { + padding: 3rem; + border: 2px dashed rgba(43, 43, 43, 0.2); + text-align: center; + cursor: pointer; + position: relative; + } + + .file-upload-zone:hover { + border-color: var(--rust); + background: rgba(193, 80, 61, 0.03); + } + + .file-upload-zone.drag-over { + border-color: var(--sage); + background: rgba(139, 154, 126, 0.1); + transform: scale(1.02); + } + + .upload-icon { + font-size: 3rem; + margin-bottom: 1rem; + opacity: 0.3; + } + + .upload-text { + font-size: 1rem; + margin-bottom: 0.5rem; + } + + .upload-hint { + font-size: 0.875rem; + opacity: 0.6; + } + + input[type="file"] { + display: none; + } + + .selected-file { + margin-top: 1.5rem; + padding: 1.5rem; + background: var(--cream); + border-left: 4px solid var(--sage); + } + + .selected-file-name { + font-weight: 600; + margin-bottom: 0.5rem; + } + + .selected-file-size { + font-size: 0.875rem; + opacity: 0.7; + } + + .submit-button { + width: 100%; + margin-top: 2rem; + padding: 1.25rem 2rem; + background: var(--charcoal); + color: var(--cream); + border: none; + font-size: 1rem; + font-weight: 600; + letter-spacing: 0.1em; + text-transform: uppercase; + cursor: pointer; + transition: all 0.3s ease; + position: relative; + overflow: hidden; + } + + .submit-button::before { + content: ''; + position: absolute; + top: 50%; + left: 50%; + width: 0; + height: 0; + background: var(--rust); + border-radius: 50%; + transform: translate(-50%, -50%); + transition: width 0.6s ease, height 0.6s ease; + } + + .submit-button:hover::before { + width: 300%; + height: 300%; + } + + .submit-button:hover { + color: white; + } + + .submit-button span { + position: relative; + z-index: 1; + } + + .submit-button:disabled { + opacity: 0.4; + cursor: not-allowed; + } + + .message { + margin-top: 2rem; + padding: 1.5rem; + border-left: 4px solid var(--sage); + background: rgba(139, 154, 126, 0.1); + animation: fadeIn 0.5s ease forwards; + } + + .message.error { + border-left-color: var(--rust); + background: rgba(193, 80, 61, 0.1); + } + + .message.success { + border-left-color: var(--sage); + background: rgba(139, 154, 126, 0.1); + } + + .loading { + display: inline-block; + width: 20px; + height: 20px; + border: 2px solid rgba(43, 43, 43, 0.1); + border-radius: 50%; + border-top-color: var(--charcoal); + animation: spin 1s linear infinite; + } + + @keyframes spin { + to { + transform: rotate(360deg); + } + } + + @media (max-width: 768px) { + main { + grid-template-columns: 1fr; + } + + h1 { + font-size: 3rem; + } + } + </style> +</head> +<body> + <div class="container"> + <header> + <div class="subtitle">LibEuFin EbiSync</div> + <h1>File Submission Portal</h1> + <div class="version-text" id="versionText">Initializing...</div> + </header> + + <main> + <section> + <h2>Choose Order</h2> + <div class="orders-list" id="ordersList"> + <div style="text-align: center; padding: 2rem; opacity: 0.5;"> + <div class="loading"></div> + <p style="margin-top: 1rem;">Loading orders...</p> + </div> + </div> + </section> + + <section> + <h2>Submit File</h2> + <div class="file-upload-zone" id="uploadZone"> + <div class="upload-icon">📄</div> + <div class="upload-text">Drop XML file here or click to browse</div> + <div class="upload-hint">Accepts .xml files only</div> + <input type="file" id="fileInput" accept=".xml,application/xml,text/xml"> + </div> + <div id="selectedFileInfo"></div> + <button class="submit-button" id="submitButton" disabled> + <span>Submit Document</span> + </button> + <div id="messageArea"></div> + </section> + </main> + </div> + + <script> + let selectedOrder = null; + let selectedFile = null; + + // Initialize + function init() { + loadConfig() + loadOrders(); + setupEventListeners(); + } + + async function loadConfig() { + try { + const response = await fetch('/config'); + const data = await response.json(); + document.getElementById('versionText').textContent = `${data.spa_version} (${data.version})`; + } catch (error) { + document.getElementById('versionText').textContent = 'Error'; + showMessage('Unable to connect to backend server', 'error'); + } + } + + async function loadOrders() { + try { + const response = await fetch('/submit'); + const data = await response.json(); + + const ordersList = document.getElementById('ordersList'); + ordersList.innerHTML = ''; + + data.orders.forEach(order => { + const orderCard = document.createElement('div'); + orderCard.className = 'order-card'; + orderCard.innerHTML = ` + <div class="order-id">${order.id}</div> + <div class="order-description">${order.description}</div> + `; + orderCard.onclick = () => selectOrder(order, orderCard); + ordersList.appendChild(orderCard); + }); + } catch (error) { + document.getElementById('ordersList').innerHTML = + '<p style="text-align: center; opacity: 0.5;">Failed to load orders</p>'; + } + } + + function selectOrder(order, element) { + document.querySelectorAll('.order-card').forEach(card => { + card.classList.remove('selected'); + }); + element.classList.add('selected'); + selectedOrder = order; + updateSubmitButton(); + } + + function setupEventListeners() { + const uploadZone = document.getElementById('uploadZone'); + const fileInput = document.getElementById('fileInput'); + + uploadZone.addEventListener('click', () => fileInput.click()); + + uploadZone.addEventListener('dragover', (e) => { + e.preventDefault(); + uploadZone.classList.add('drag-over'); + }); + + uploadZone.addEventListener('dragleave', () => { + uploadZone.classList.remove('drag-over'); + }); + + uploadZone.addEventListener('drop', (e) => { + e.preventDefault(); + uploadZone.classList.remove('drag-over'); + const files = e.dataTransfer.files; + if (files.length > 0) { + handleFileSelect(files[0]); + } + }); + + fileInput.addEventListener('change', (e) => { + if (e.target.files.length > 0) { + handleFileSelect(e.target.files[0]); + } + }); + + document.getElementById('submitButton').addEventListener('click', submitFile); + } + + function handleFileSelect(file) { + if (!file.name.endsWith('.xml')) { + showMessage('Please select an XML file', 'error'); + return; + } + + selectedFile = file; + const infoDiv = document.getElementById('selectedFileInfo'); + infoDiv.innerHTML = ` + <div class="selected-file"> + <div class="selected-file-name">📎 ${file.name}</div> + <div class="selected-file-size">${(file.size / 1024).toFixed(2)} KB</div> + </div> + `; + updateSubmitButton(); + } + + function updateSubmitButton() { + const button = document.getElementById('submitButton'); + button.disabled = !(selectedOrder && selectedFile); + } + + async function submitFile() { + const button = document.getElementById('submitButton'); + button.disabled = true; + button.innerHTML = '<span><div class="loading" style="display: inline-block; vertical-align: middle; margin-right: 10px;"></div>Submitting...</span>'; + + const formData = new FormData(); + formData.append('order', selectedOrder.id); + formData.append('file', selectedFile); + + clearMessage() + + try { + const response = await fetch('/submit', { + method: 'POST', + body: formData + }); + const data = await response.json(); + if (response.ok) { + showMessage(`Successfully submitted ${selectedFile.name} with order ${selectedOrder.id} as ${data.order}`, 'success'); + // Reset form + selectedFile = null; + document.getElementById('selectedFileInfo').innerHTML = ''; + document.getElementById('fileInput').value = ''; + } else { + showMessage(`${data.code} - ${data.hint ?? 'Submission failed'}`, 'error'); + } + } catch (error) { + showMessage('Network error: ' + error.message, 'error'); + } finally { + button.disabled = false; + button.innerHTML = '<span>Submit Document</span>'; + updateSubmitButton(); + } + } + + function showMessage(text, type = 'info') { + const messageArea = document.getElementById('messageArea'); + messageArea.innerHTML = ` + <div class="message ${type}"> + ${text} + </div> + `; + } + + function clearMessage() { + const messageArea = document.getElementById('messageArea'); + messageArea.innerHTML = ''; + } + + // Start application + init(); + </script> +</body> +</html> +\ No newline at end of file diff --git a/libeufin-ebisync/src/api.rs b/libeufin-ebisync/src/api.rs @@ -0,0 +1,220 @@ +/* +* 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::sync::Arc; + +use axum::{ + Json, Router, + body::Bytes, + extract::{Multipart, State}, + response::Redirect, + routing::get, +}; +use compact_str::{CompactString, ToCompactString}; +use libeufin_ebics::{ + cli::EbicsLogs, + ebics::{EbicsClient, EbicsErrKind, administrative::OrderInfo}, + keys::{BankKeys, ClientKeys}, +}; +use reqwest::StatusCode; +use serde::{Deserialize, Serialize}; +use sqlx::PgPool; +use taler_api::{ + api::{RouterUtils, TalerRouter as _}, + auth::AuthMethod, + error::{ApiResult, bad_request, failure, failure_status}, +}; +use taler_build::long_version; +use taler_common::{api::LibtoolVersion, error_code::ErrorCode}; +use taler_macros::api_config; +use tower_http::services::ServeDir; + +use crate::config::EbisyncHostCfg; + +pub struct EbisyncState { + pub db: PgPool, + pub cfg: EbisyncHostCfg, + pub client: ClientKeys, + pub bank: BankKeys, +} + +impl EbisyncState { + pub fn client<'a>(&'a self) -> anyhow::Result<EbicsClient<'a>> { + EbicsClient::new(self.cfg.ebics(), EbicsLogs { dir: None }) + } + + pub async fn orders(&self) -> anyhow::Result<Vec<OrderInfo>> { + let mut hkd = self + .client()? + .hkd(&self.db, &self.client, &self.bank, false) + .await?; + hkd.partner.orders.retain(|it| it.order.is_upload()); + Ok(hkd.partner.orders) + } +} + +#[api_config("taler-ebisync")] +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +pub struct Config<'a> { + spa_version: &'a str, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct ListSubmitOrders { + pub orders: Vec<SubmitOrder>, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct SubmitOrder { + pub id: String, + pub description: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct SyncSubmit { + pub order: CompactString, +} + +pub fn sync_api(state: Arc<EbisyncState>, spa: &str, auth: AuthMethod) -> Router { + Router::new() + .route( + "/config", + get(async || { + Json(Config { + name: (), + version: LibtoolVersion::new(0, 0, 0), + implementation: Some("urn:net:taler:specs:libeufin-ebbisync:taler-rust"), + spa_version: long_version(), + }) + }), + ) + .nest_service("/webui/", ServeDir::new(spa)) + .route("/", get(async || Redirect::permanent("/webui/"))) + .route( + "/submit", + get(async |State(state): State<Arc<EbisyncState>>| { + ApiResult::Ok(Json(ListSubmitOrders { + orders: state + .orders() + .await + .map_err(|e| { + failure_status( + ErrorCode::END, + e.to_compact_string(), + StatusCode::BAD_GATEWAY, + ) + })? + .into_iter() + .map(|it| SubmitOrder { + id: it.order.to_string(), + description: it.description, + }) + .collect(), + })) + }) + .post( + async |State(state): State<Arc<EbisyncState>>, mut multipart: Multipart| { + let mut order_id: Option<String> = None; + let mut xml: Option<Bytes> = None; + + while let Some(field) = multipart + .next_field() + .await + .map_err(|_| bad_request("Invalid multipart"))? + { + let name = field.name().unwrap_or("").to_string(); + + match name.as_str() { + "order" => { + order_id = Some( + field + .text() + .await + .map_err(|_| bad_request("Invalid order field"))?, + ); + } + _ => { + // treat as file + let data = field + .bytes() + .await + .map_err(|_| bad_request("Invalid file"))?; + + xml = Some(data); + } + } + } + + let order_id = order_id.ok_or_else(|| { + failure(ErrorCode::GENERIC_PARAMETER_MISSING, "Missing order ID") + .with_path("order") + })?; + + let xml = xml.ok_or_else(|| { + failure(ErrorCode::GENERIC_PARAMETER_MISSING, "Missing file") + })?; + + // --- match lookup --- + let match_order = state + .orders() + .await + .map_err(|e| { + failure_status( + ErrorCode::END, + e.to_compact_string(), + StatusCode::BAD_GATEWAY, + ) + })? + .into_iter() + .find(|o| o.order.to_string() == order_id.as_str()) + .ok_or_else(|| { + failure_status( + ErrorCode::END, + format!("Unknown order '{}'", order_id), + StatusCode::NOT_FOUND, + ) + })?; + + // --- upload with error mapping --- + let order = state + .client() + .map_err(|e| { + failure_status( + ErrorCode::END, + e.to_compact_string(), + StatusCode::BAD_GATEWAY, + ) + })? + .upload(&state.client, &state.bank, &match_order.order, xml.as_ref()) + .await + .map_err(|e| match e.kind { + EbicsErrKind::Code { .. } => { + failure_status(ErrorCode::END, e, StatusCode::CONFLICT) + } + _ => failure_status(ErrorCode::END, e, StatusCode::BAD_GATEWAY), + })?; + + ApiResult::Ok(Json(SyncSubmit { order })) + }, + ), + ) + .with_state(state) + .auth(auth, "realm") + .finalize() +} diff --git a/libeufin-ebisync/src/azure.rs b/libeufin-ebisync/src/azure.rs @@ -0,0 +1,219 @@ +/* +* This file is part of LibEuFin. +* Copyright (C) 2026 Taler Systems S.A. + +* LibEuFin is free software; you can redistribute it and/or modify +* it under the terms of the GNU Affero General Public License as +* published by the Free Software Foundation; either version 3, or +* (at your option) any later version. + +* LibEuFin is distributed in the hope that it will be useful, but +* WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY +* or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Affero General +* Public License for more details. + +* You should have received a copy of the GNU Affero General Public +* License along with LibEuFin; see the file COPYING. If not, see +* <http://www.gnu.org/licenses/> +*/ + +use std::str::FromStr; + +use aws_lc_rs::hmac::{self, HMAC_SHA256}; +use axum::{ + body::Bytes, + http::{HeaderName, HeaderValue}, +}; +use compact_str::{CompactString, ToCompactString}; +use http_client::{ApiErr, Client, ClientErr, builder::Req}; +use jiff::{Timestamp, Zoned, tz::TimeZone}; +use reqwest::{ + Method, StatusCode, Url, + header::{ + AUTHORIZATION, CONTENT_ENCODING, CONTENT_LANGUAGE, CONTENT_TYPE, IF_MATCH, + IF_MODIFIED_SINCE, IF_NONE_MATCH, IF_RANGE, IF_UNMODIFIED_SINCE, + }, +}; +use taler_common::encoding::base64; +use thiserror::Error; + +const API_VERSION: &str = "2025-11-05"; + +#[derive(Debug, Error)] + +pub enum AzureErr { + #[error("{status} {code}")] + Api { + status: StatusCode, + code: CompactString, + }, + #[error(transparent)] + Client(#[from] ClientErr), +} + +pub type AzureResult = Result<(), ApiErr<AzureErr>>; + +pub struct AzureBlobStorage<'a> { + base_url: Url, + account: &'a str, + key: hmac::Key, + client: &'a Client, +} + +impl<'a> AzureBlobStorage<'a> { + pub fn new( + base_url: &'a str, + account: &'a str, + key: &str, + client: &'a Client, + ) -> anyhow::Result<Self> { + let decoded = base64::decode(key)?; + Ok(Self { + base_url: Url::from_str(base_url)?, + account, + key: hmac::Key::new(HMAC_SHA256, &decoded), + client, + }) + } + + async fn req(&self, req: Req) -> AzureResult { + // Set required headers (x-ms-date and x-ms-version) + // Azure uses x-ms-date instead of the standard Date header for signing + let date = Zoned::new(Timestamp::now(), TimeZone::UTC) + .strftime("%a, %d %b %Y %H:%M:%S GMT") + .to_string(); + let req = req + .header("x-ms-date", date) + .header("x-ms-version", API_VERSION); + + // Calculate the HMAC-SHA256 signature + let signature = { + let mut ctx = aws_lc_rs::hmac::Context::with_key(&self.key); + + let add_str = |sign: &mut aws_lc_rs::hmac::Context, value: &str| { + sign.update(value.as_bytes()); + sign.update(b"\n"); + }; + let add_header = |sign: &mut aws_lc_rs::hmac::Context, name: &HeaderName| { + sign.update( + req.headers() + .get(name) + .map(|it| it.as_bytes()) + .unwrap_or_default(), + ); + sign.update(b"\n"); + }; + // 1. VERB + add_str(&mut ctx, req.method().as_str()); + // 2. Content-Encoding + add_header(&mut ctx, &CONTENT_ENCODING); + // 3. Content-Language + add_header(&mut ctx, &CONTENT_LANGUAGE); + // 4. Content-Length (empty string if zero for modern versions) + let len = req.body().len(); + if len != 0 { + add_str(&mut ctx, &len.to_compact_string()); + } else { + add_str(&mut ctx, ""); + } + // 5. Content-MD5 + add_header(&mut ctx, &HeaderName::from_static("content-md5")); + // 6. Content-Type + add_header(&mut ctx, &CONTENT_TYPE); + // 7. Date + add_str(&mut ctx, ""); // Must be empty as x-ms-date is used) + // 8. If-Modified-Since + add_header(&mut ctx, &IF_MODIFIED_SINCE); + // 9. If-Match + add_header(&mut ctx, &IF_MATCH); + // 10. If-None-Match + add_header(&mut ctx, &IF_NONE_MATCH); + // 11. If-Unmodified-Since + add_header(&mut ctx, &IF_UNMODIFIED_SINCE); + // 12. Range + add_header(&mut ctx, &IF_RANGE); + // 13. CanonicalizedHeaders + // This includes all x-ms- headers, converted to lowercase, sorted, and concatenated. + let mut headers: Vec<_> = req + .headers() + .iter() + .filter(|(n, _)| n.as_str().starts_with("x-ms-")) + .collect(); + headers.sort_unstable_by_key(|(n, _)| n.as_str()); + // TODO should we group them ? + for (n, v) in headers { + ctx.update(n.as_str().as_bytes()); + ctx.update(b":"); + ctx.update(v.as_bytes()); // TODO replace linear whitesoace ? + ctx.update(b"\n"); + } + + // 14. CanonicalizedResource + // This includes the account name, the path, and canonicalized query parameters. + ctx.update(b"/"); + ctx.update(self.account.as_bytes()); + ctx.update(req.url().path().trim_end_matches('/').as_bytes()); + let mut params: Vec<_> = req.url().query_pairs().collect(); + params.sort_unstable(); + for (key, value) in params { + ctx.update(b"\n"); + ctx.update(key.as_bytes()); + ctx.update(b":"); + ctx.update(value.as_bytes()); + } + ctx.sign() + }; + + // Add the Authorization header + let req = req.header( + AUTHORIZATION, + format!("SharedKey {}:{}", self.account, base64::fmt(signature)), + ); + + // Send it and handle error + let (ctx, res) = req.send().await.map_err(|(ctx, e)| ctx.wrap(e.into()))?; + if !res.status().is_success() { + Err(ctx.wrap(AzureErr::Api { + status: res.status(), + code: res + .headers() + .get(HeaderName::from_static("x-ms-error-code")) + .and_then(|it| it.to_str().ok()) + .unwrap_or_default() + .into(), + })) + } else { + Ok(()) + } + } + + pub async fn create_container(&self, name: &'static str) -> AzureResult { + self.req( + Req::new(self.client, Method::PUT, &self.base_url, name).query("restype", "container"), + ) + .await + } + + pub async fn put_blob( + &self, + container: &str, + name: &str, + content: Bytes, + ty: HeaderValue, + ) -> AzureResult { + self.req( + Req::new( + self.client, + Method::PUT, + &self.base_url, + format!("{container}/{name}"), + ) + .content(content, ty) + .header( + HeaderName::from_static("x-ms-blob-type"), + HeaderValue::from_static("BlockBlob"), + ), + ) + .await + } +} diff --git a/libeufin-ebisync/src/config.rs b/libeufin-ebisync/src/config.rs @@ -0,0 +1,250 @@ +/* +* 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::{cell::OnceCell, time::Duration}; + +use compact_str::CompactString; +use jiff::civil::Time; +use libeufin_ebics::{ + config::{EbicsHostCfg, EbicsKeysCfg, EbicsSetupCfg}, + keys::{BankKeys, ClientKeys}, +}; +use reqwest::Url; +use taler_api::{ + Serve, + config::{AuthCfg, DbCfg}, +}; +use taler_common::{ + config::{Config, Section, ValueErr, parser::ConfigSource}, + map_config, +}; + +pub const CONFIG_SOURCE: ConfigSource = + ConfigSource::new("libeufin-ebisync", "ebisync", "libeufin-ebisync"); + +pub fn parse_db_cfg(cfg: &Config) -> Result<DbCfg, ValueErr> { + DbCfg::parse(cfg.section("ebisyncdb-postgres")) +} + +pub struct EbisyncSetupCfg { + pub enc: Option<Vec<u8>>, + pub auth: Option<Vec<u8>>, +} + +impl EbisyncSetupCfg { + pub fn parse(cfg: &Config) -> Result<Self, ValueErr> { + let s = cfg.section("ebisync-setup"); + Ok(Self { + enc: s.hex("bank_encryption_pub_key_hash").opt()?, + auth: s.hex("bank_authentication_pub_key_hash").opt()?, + }) + } +} + +pub enum Destination { + None, + AzureBlobStorage { + api: Url, + name: String, + key: String, + container: String, + }, +} + +pub enum Source { + None, + SyncAPI(AuthCfg), +} + +pub struct EbisyncFetchCfg { + pub frequency: Duration, + pub frequency_raw: String, + pub checkpoint_time: Time, + pub destination: Destination, +} + +impl EbisyncFetchCfg { + pub fn parse(cfg: &Config) -> Result<Self, ValueErr> { + let s = cfg.section("ebisync-fetch"); + Ok(Self { + frequency: s.duration("frequency").require()?, + frequency_raw: s.str("frequency").require()?, + checkpoint_time: s.time("checkpoint_time_of_day").require()?, + destination: map_config!(s, "ebics file destination", "destination", + "none" => { Destination::None }, + "azure-blob-storage" => { + Destination::AzureBlobStorage { api: s.base_url("azure_api_url").require()?, name: s.str("azure_account_name").require()?, key: s.str("azure_account_key").require()?, container: s.str("azure_container").require()? } + } + ).require()? + }) + } +} + +pub struct EbisyncSubmitCfg { + pub source: Source, +} + +impl EbisyncSubmitCfg { + pub fn parse(cfg: &Config) -> Result<Self, ValueErr> { + let s = cfg.section("ebisync-fetch"); + Ok(Self { + source: map_config!(s, "ebics file destination", "destination", + "none" => { Source::None }, + "ebisync-api" => { Source::SyncAPI(AuthCfg::parse(&s)?) } + ) + .require()?, + }) + } +} + +pub struct EbisyncServeCfg { + pub serve: Serve, + pub spa_path: String, +} + +impl EbisyncServeCfg { + pub fn parse(cfg: &Config) -> Result<Self, ValueErr> { + let s = cfg.section("ebisync-httpd"); + Ok(Self { + serve: Serve::parse(&s)?, + spa_path: s.path("spa").require()?, + }) + } +} + +#[derive(Clone)] +pub struct EbisyncHostCfg { + pub base_url: Url, + pub unix_path: Option<String>, + pub host_id: CompactString, + pub user_id: CompactString, + pub partner_id: CompactString, +} + +impl EbisyncHostCfg { + pub fn parse(s: &Section) -> Result<Self, ValueErr> { + Ok(Self { + base_url: s.url("host_base_url").require()?, + unix_path: s.path("UNIXPATH").opt()?, + host_id: s.cstr("host_id").require()?, + user_id: s.cstr("user_id").require()?, + partner_id: s.cstr("partner_id").require()?, + }) + } + + pub fn ebics<'a>(&'a self) -> EbicsHostCfg<'a> { + EbicsHostCfg { + base_url: self.base_url.as_str(), + unix_path: self.unix_path.as_deref(), + host_id: &self.host_id, + user_id: &self.user_id, + partner_id: &self.partner_id, + } + } +} + +pub struct EbisyncCfg { + pub cfg: Config, + pub bank: String, + pub client: String, + pub host: EbisyncHostCfg, + pub fetch: OnceCell<EbisyncFetchCfg>, + pub submit: OnceCell<EbisyncSubmitCfg>, + pub setup: OnceCell<EbisyncSetupCfg>, + pub serve: OnceCell<EbisyncServeCfg>, + pub db_cfg: DbCfg, +} + +impl EbisyncCfg { + pub fn parse(cfg: Config) -> Result<Self, ValueErr> { + let s = cfg.section("ebisync"); + Ok(Self { + host: EbisyncHostCfg::parse(&s)?, + bank: s.path("bank_public_keys_file").require()?, + client: s.path("client_private_keys_file").require()?, + fetch: OnceCell::new(), + submit: OnceCell::new(), + setup: OnceCell::new(), + serve: OnceCell::new(), + db_cfg: parse_db_cfg(&cfg)?, + cfg, + }) + } + + pub fn fetch(&self) -> Result<&EbisyncFetchCfg, ValueErr> { + // TODO use get_or_try_init when stable + if let Some(fetch) = self.fetch.get() { + return Ok(fetch); + } + let fetch = EbisyncFetchCfg::parse(&self.cfg)?; + self.fetch.set(fetch).ok(); + Ok(self.fetch.get().unwrap()) + } + + pub fn submit(&self) -> Result<&EbisyncSubmitCfg, ValueErr> { + // TODO use get_or_try_init when stable + if let Some(submit) = self.submit.get() { + return Ok(submit); + } + let submit = EbisyncSubmitCfg::parse(&self.cfg)?; + self.submit.set(submit).ok(); + Ok(self.submit.get().unwrap()) + } + + pub fn serve(&self) -> Result<&EbisyncServeCfg, ValueErr> { + // TODO use get_or_try_init when stable + if let Some(serve) = self.serve.get() { + return Ok(serve); + } + let serve = EbisyncServeCfg::parse(&self.cfg)?; + self.serve.set(serve).ok(); + Ok(self.serve.get().unwrap()) + } + + pub fn ebics_setup<'a>(&'a self) -> Result<EbicsSetupCfg<'a>, ValueErr> { + let setup = { + if let Some(setup) = self.setup.get() { + setup + } else { + let setup = EbisyncSetupCfg::parse(&self.cfg)?; + self.setup.set(setup).ok(); + self.setup.get().unwrap() + } + }; + Ok(EbicsSetupCfg { + keys: EbicsKeysCfg { + bank: &self.bank, + client: &self.client, + }, + host: self.host.ebics(), + enc: setup.enc.as_deref(), + auth: setup.auth.as_deref(), + }) + } + + pub fn expect_full_keys(&self) -> anyhow::Result<(ClientKeys, BankKeys)> { + libeufin_ebics::keys::expect_full_keys( + &EbicsKeysCfg { + bank: &self.bank, + client: &self.client, + }, + "libeufin-ebisync setup", + ) + } +} diff --git a/libeufin-ebisync/src/db.rs b/libeufin-ebisync/src/db.rs @@ -0,0 +1,60 @@ +/* +* 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 sqlx::PgPool; +use taler_api::config::DbCfg; + +use crate::config::CONFIG_SOURCE; + +pub async fn pool(cfg: &DbCfg) -> anyhow::Result<PgPool> { + let pool = taler_common::db::pool(cfg.cfg.clone(), CONFIG_SOURCE.component_name).await?; + Ok(pool) +} + +pub async fn dbinit(cfg: &DbCfg, reset: bool) -> anyhow::Result<PgPool> { + let pool = taler_common::db::pool(cfg.cfg.clone(), CONFIG_SOURCE.component_name).await?; + let mut db = pool.acquire().await?; + taler_common::db::dbinit( + &mut db, + cfg.sql_dir.as_ref(), + CONFIG_SOURCE.component_name, + reset, + ) + .await?; + Ok(pool) +} + +#[cfg(test)] +pub mod test { + + use libeufin_ebics::db::test::ebics_routine; + use sqlx::{PgPool, Postgres, pool::PoolConnection}; + + use crate::config::CONFIG_SOURCE; + + pub async fn db_setup() -> (PoolConnection<Postgres>, PgPool) { + taler_test_utils::db::db_test_setup(CONFIG_SOURCE).await + } + + #[tokio::test] + pub async fn ebics() { + let (_, db) = db_setup().await; + ebics_routine(&db).await; + } +} diff --git a/libeufin-ebisync/src/lib.rs b/libeufin-ebisync/src/lib.rs @@ -0,0 +1,562 @@ +/* +* 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/> +*/ + +#![allow(clippy::too_many_arguments)] + +use std::{ + fmt::Write as _, + io::{Cursor, Read as _}, + sync::Arc, + time::Duration, +}; + +use anyhow::anyhow; +use axum::{body::Bytes, http::HeaderValue}; +use compact_str::CompactStringExt; +use http_client::Client; +use jiff::{Timestamp, Zoned, civil::Date, tz::TimeZone}; +use libeufin_ebics::{ + cli::{EbicsArgs, EbicsLogs}, + db::{get_task_status, update_task_status}, + ebics::{ + EbicsClient, EbicsErrKind, + administrative::{AccountInfo, HKD, OrderInfo}, + ebics_code::EbicsReturnCode, + order::{Order, OrderDoc}, + }, + keys::{BankKeys, ClientKeys}, + ws::listen_for_notification, +}; +use sqlx::PgPool; +use taler_api::api::TalerRouter; +use taler_build::long_version; +use taler_common::{CommonArgs, cli::ConfigCmd, config::Config, types::utils::date_to_utc_ts}; +use tokio::{time::timeout, try_join}; +use tracing::{debug, error, info, trace}; + +use crate::{ + api::{EbisyncState, sync_api}, + azure::AzureBlobStorage, + config::{Destination, EbisyncCfg, Source, parse_db_cfg}, + db::{dbinit, pool}, +}; + +pub mod api; +pub mod azure; +pub mod config; +pub mod db; + +// KV +const CHECKPOINT_KEY: &str = "checkpoint"; +const FETCH_TASK_KEY: &str = "fetch_task"; + +#[derive(clap::Subcommand, Debug)] +pub enum Cmd { + /// Initialize libeufin-ebisync database + Dbinit { + /// Reset database (DANGEROUS: All existing data is lost) + #[arg(short, long)] + reset: bool, + }, + /// Set up the EBICS subscriber + Setup { + #[command(flatten)] + ebics_logs: EbicsLogs, + + /// Resubmits all the keys to the bank + #[arg(long)] + force_keys_resubmission: bool, + + /// Accepts the bank keys without interactively asking the user + #[arg(long)] + auto_accept_keys: bool, + + /// Generates the PDF with the client public keys to send to the bank + #[arg(long)] + generate_registration_pdf: bool, + }, + /// Downloads EBICS files from the bank and store them in the configured destination + Fetch { + #[clap(flatten)] + ebics: EbicsArgs, + + /// Only supported in --transient mode, this option lets specify the earliest timestamp of the downloaded documents + #[arg(long, value_name = "YYYY-MM-DD")] + pinned_start: Option<Date>, + + /// Only supported in --transient mode, do not consume fetched documents + #[arg(long, requires = "transient")] + peek: bool, + + /// Only supported in --transient mode, run a checkpoint + #[arg(long, requires = "transient")] + checkpoint: bool, + /// Check whether a destination is configured. Exit with 0 if at destination is configured, otherwise 1 + #[arg(long)] + check: bool, + }, + /// Run libeufin-ebisync HTTP server + Serve { + /// Check whether an API is in use (if it's useful to start the HTTP server). Exit with 0 if at least one API is enabled, otherwise 1 + #[arg(long)] + check: bool, + }, + #[command(subcommand)] + Config(ConfigCmd), +} + +#[derive(clap::Parser, Debug)] +#[command(long_version = long_version(), about, long_about = None)] +pub struct Args { + #[clap(flatten)] + pub common: CommonArgs, + + #[command(subcommand)] + pub cmd: Cmd, +} + +pub async fn ebics_setup( + ebics: &EbicsClient<'_>, + cfg: &EbisyncCfg, + db: &PgPool, + force_keys_resubmission: bool, + generate_registration_pdf: bool, + auto_accept_keys: bool, +) -> anyhow::Result<()> { + let (client, bank) = libeufin_ebics::setup::ebics_setup( + ebics, + &cfg.ebics_setup()?, + force_keys_resubmission, + generate_registration_pdf, + auto_accept_keys, + ) + .await?; + + // Check account information + info!(target: "setup", "Doing administrative request HKD"); + let HKD { partner, .. } = ebics.hkd(db, &client, &bank, false).await?; + // Debug logging + let fmt = std::fmt::from_fn(|f| { + if partner.name.is_some() || !partner.accounts.is_empty() { + f.write_str("Partner Info: ")?; + if let Some(name) = &partner.name { + write!(f, "'{name}'")?; + } + for AccountInfo { + currency, + iban, + bic, + } in &partner.accounts + { + write!(f, "{currency}-{iban}{bic}")?; + } + f.write_char('\n')?; + } + writeln!(f, "Supported orders:")?; + for OrderInfo { order, description } in &partner.orders { + writeln!(f, "- {order}: {description}")?; + } + Ok(()) + }); + debug!(target: "setup", "{fmt}"); + info!(target: "setup", "EBICS ready"); + + eprintln!("setup ready"); + Ok(()) +} + +pub enum DestinationClient<'a> { + AzureBlobStorage { + client: AzureBlobStorage<'a>, + container: &'a str, + }, +} + +impl<'a> DestinationClient<'a> { + pub fn prepare(dest: &'a Destination, client: &'a Client) -> anyhow::Result<Option<Self>> { + Ok(match dest { + Destination::None => None, + Destination::AzureBlobStorage { + api, + name: account, + key, + container, + } => Some(DestinationClient::AzureBlobStorage { + client: AzureBlobStorage::new(api.as_str(), account, key, client)?, + container, + }), + }) + } + + pub async fn upload(&self, name: &str, xml: Bytes) -> anyhow::Result<()> { + match self { + DestinationClient::AzureBlobStorage { client, container } => { + client + .put_blob( + container, + name, + xml, + HeaderValue::from_static("application/xml"), + ) + .await?; + } + } + Ok(()) + } +} + +pub async fn ebics_fetch( + ebics: &EbicsClient<'_>, + cfg: &EbisyncCfg, + client: &ClientKeys, + bank: &BankKeys, + db: &PgPool, + pinned_start: &Option<Timestamp>, + peek: bool, + check: bool, + transient: bool, + transient_checkpoint: bool, +) -> anyhow::Result<()> { + let fetch_cfg = cfg.fetch()?; + let (sender, mut receiver) = tokio::sync::mpsc::channel::<Vec<Order>>(10); + let http = http_client::client(); + let dest = DestinationClient::prepare(&fetch_cfg.destination, &http)?; + + if check { + if dest.is_none() { + info!("No destination configured, not starting the fetcher"); + std::process::exit(1); + } + } else if let Some(dest) = dest { + let upload = async |orders: &[Order], since: Option<Timestamp>| -> anyhow::Result<()> { + for order in orders { + if let Err(e) = ebics + .download( + db, + client, + bank, + order, + &since.map(|it| (it, Timestamp::now())), + transient && peek, + async |content| { + if order.doc() == Some(OrderDoc::acknowledgement) { + // TODO HAC + } else { + let mut z = zip::ZipArchive::new(Cursor::new(content))?; + for i in 0..z.len() { + let mut file = z.by_index(i)?; + trace!(target: "fetch", "upload {}", file.name()); + let mut buf = Vec::new(); + file.read_to_end(&mut buf)?; + dest.upload(file.name(), buf.into()) + .await + .map_err(|e| EbicsErrKind::Custom(e.to_string().into()))?; + } + } + Ok(()) + }, + ) + .await + { + if let EbicsErrKind::Code { + bank: EbicsReturnCode::EBICS_NO_DOWNLOAD_DATA_AVAILABLE, + .. + } = e.kind + { + continue; + } + return Err(e.into()); + } + } + Ok(()) + }; + let fetch = async { + if transient { + info!(target: "fetch", "Transient mode: fetching once and returning"); + } else { + info!(target: "fetch", "Running with a frequency of {}", fetch_cfg.frequency_raw); + } + + let mut last_fetch = Timestamp::UNIX_EPOCH; + loop { + let now = Timestamp::now(); + let checkpoint = get_task_status(db, CHECKPOINT_KEY) + .await? + .unwrap_or_default(); + let next_fetch = last_fetch + fetch_cfg.frequency; + let next_checkpoint = { + if let Some(last_trial) = checkpoint.last_trial { + // We run today at checkpoint_time + let checkpoint_date = Zoned::new(now, TimeZone::UTC) + .with() + .time(fetch_cfg.checkpoint_time) + .build() + .unwrap(); + // If we already ran today we ran tomorrow + if last_trial > checkpoint_date.timestamp() { + checkpoint_date.tomorrow().unwrap().timestamp() + } else { + checkpoint_date.timestamp() + } + } else { + // We never ran, we must checkpoint now + now + } + }; + + let mut success = true; + if + // Run transient checkpoint at request + (transient && transient_checkpoint) + // Or run recurrent checkpoint + || (!transient && now > next_checkpoint) + { + info!(target: "fetch", "Running checkpoint"); + + let since = if let Some(pinned_start) = pinned_start + && transient + && checkpoint + .last_successfull + .map(|it| *pinned_start <= it) + .unwrap_or(true) + { + Some(*pinned_start) + } else { + checkpoint.last_successfull + }; + let res = async { + // We fetch HKD to only fetch supported EBICS orders and get the document versions + let hkd = ebics.hkd(db, client, bank, false).await?; + let mut supported_orders = hkd + .partner + .orders + .into_iter() + .map(|it| it.order) + .collect::<Vec<_>>(); + debug!(target: "fetch", "HKD: {}", supported_orders.iter().map(|it| it.to_string()).join_compact(", ")); + supported_orders + .retain(|it| it.is_downloadable()); + upload( &supported_orders, since).await + } + .await; + if let Err(e) = res { + success = false; + error!(target: "fetch", "{e}"); + } + try_join!( + update_task_status(db, CHECKPOINT_KEY, &now, success), + update_task_status(db, FETCH_TASK_KEY, &now, success) + )?; + last_fetch = now; + } else if transient || now > next_fetch { + if !transient { + info!(target: "fetch", "Running at frequency"); + } + let res = async { + // We fetch HAA to only fetch pending & supported EBICS orders and get the document versions + let mut haa = ebics.haa(db, client, bank, false).await?; + debug!(target: "fetch", "HAA: {}", haa.orders.iter().map(|it| it.to_string()).join_compact(", ")); + haa.orders + .retain(|it| it.is_downloadable()); + upload( &haa.orders, *pinned_start).await + } + .await; + if let Err(e) = res { + success = false; + error!(target: "fetch", "{e}"); + } + update_task_status(db, FETCH_TASK_KEY, &now, success).await?; + last_fetch = now; + } + + if transient { + if success { + return anyhow::Ok(()); + } else { + return Err(anyhow!("fetch failed")); + } + } + + let delay = now.duration_until(next_fetch.min(next_checkpoint)); + let tx = timeout( + Duration::from_millis(delay.abs().as_millis() as u64), + receiver.recv(), + ) + .await; + if let Ok(Some(notification)) = tx { + info!(target: "fetch", "Running at real-time notifications reception"); + upload(&notification, None).await?; + } + } + }; + if transient { + fetch.await?; + } else { + tokio::try_join!(fetch, async { + listen_for_notification(ebics, db, client, bank, sender).await; + Ok(()) + })?; + } + } + + Ok(()) +} + +pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { + match cmd { + Cmd::Dbinit { reset } => { + let cfg = parse_db_cfg(&cfg)?; + dbinit(&cfg, reset).await?; + } + Cmd::Setup { + ebics_logs, + force_keys_resubmission, + auto_accept_keys, + generate_registration_pdf, + } => { + let cfg = EbisyncCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; + let ebics = EbicsClient::new(cfg.host.ebics(), ebics_logs)?; + ebics_setup( + &ebics, + &cfg, + &pool, + force_keys_resubmission, + generate_registration_pdf, + auto_accept_keys, + ) + .await?; + } + Cmd::Fetch { + pinned_start, + peek, + checkpoint, + check, + ebics: EbicsArgs { logs, transient }, + } => { + let cfg = EbisyncCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; + let ebics = EbicsClient::new(cfg.host.ebics(), logs)?; + let (client, bank) = cfg.expect_full_keys()?; + ebics_fetch( + &ebics, + &cfg, + &client, + &bank, + &pool, + &pinned_start.map(|it| date_to_utc_ts(&it)), + peek, + check, + transient, + transient && checkpoint, + ) + .await? + } + Cmd::Serve { check } => { + let cfg = EbisyncCfg::parse(cfg)?; + let auth = match &cfg.submit()?.source { + Source::None => None, + Source::SyncAPI(auth_cfg) => Some(auth_cfg.method()), + }; + if check { + if auth.is_none() { + info!("No source api, not starting the server"); + std::process::exit(1); + } + } else if let Some(auth) = auth { + let (client, bank) = cfg.expect_full_keys()?; + let tmp = EbisyncState { + db: pool(&cfg.db_cfg).await?, + cfg: cfg.host.clone(), + client, + bank, + }; + let cfg = cfg.serve()?; + + sync_api(Arc::new(tmp), &cfg.spa_path, auth) + .serve(&cfg.serve, None) + .await?; + } + } + Cmd::Config(cmd) => cmd.run(&cfg)?, + } + Ok(()) +} + +#[cfg(test)] +mod ebics { + use clap::Parser as _; + use libeufin_ebics::test::{EbicsState, TestBank}; + use sqlx::PgPool; + use taler_common::config::Config; + + use crate::{Args, config::CONFIG_SOURCE, db::test::db_setup, run}; + + pub async fn ebisync_cmd(cfg: &Config, cmd: &str) -> anyhow::Result<()> { + let parts = shlex::split(cmd).unwrap(); + let args = std::iter::once("libeufin-ebisync").chain(parts.iter().map(|it| it.as_str())); + + let cmd = Args::try_parse_from(args).unwrap(); + run(cfg.clone(), cmd.cmd).await + } + + async fn test_setup() -> (TestBank, Config, PgPool) { + let (_, db) = db_setup().await; + let test = TestBank::new().await; + let cfg = Config::from_mem_with_env( + CONFIG_SOURCE, + &format!( + " + [paths] + EBISYNC_HOME = {:?} + + [ebisync] + HOST_BASE_URL = http://bank.example.com/ + UNIXPATH = {} + HOST_ID = PFEBICS + USER_ID = PFC00563 + PARTNER_ID = PFC00563 + + [ebisyncdb-postgres] + CONFIG = postgresql:///{} + ", + test.dir.path(), + test.sock_path, + db.connect_options().get_database().unwrap() + ), + ) + .unwrap(); + test.sequences(&[ + EbicsState::hev, + EbicsState::ini, + EbicsState::hia, + EbicsState::hpb, + EbicsState::hkd, + EbicsState::receipt_ok, + ]); + ebisync_cmd(&cfg, "setup --auto-accept-keys").await.unwrap(); + + (test, cfg, db) + } + + #[tokio::test] + async fn setup() { + test_setup().await; + } +} diff --git a/libeufin-ebisync/src/main.rs b/libeufin-ebisync/src/main.rs @@ -0,0 +1,27 @@ +/* +* 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 clap::Parser as _; +use libeufin_ebisync::{Args, config::CONFIG_SOURCE, run}; +use taler_common::taler_main; + +fn main() { + let args = Args::parse(); + taler_main(CONFIG_SOURCE, args.common, |cfg| run(cfg, args.cmd)) +} diff --git a/libeufin-nexus/Cargo.toml b/libeufin-nexus/Cargo.toml @@ -25,10 +25,8 @@ serde.workspace = true sqlx.workspace = true compact_str.workspace = true uuid.workspace = true -shlex = "2.0" +shlex.workspace = true url = "2.5" regex = "1.12" const_format = { version = "0.2", features = ["rust_1_83"] } -zip = { version = "8.5", default-features = false, features = [ - "deflate-flate2-zlib-rs", -] } -\ No newline at end of file +zip.workspace = true +\ No newline at end of file diff --git a/libeufin-nexus/src/config.rs b/libeufin-nexus/src/config.rs @@ -19,6 +19,7 @@ use std::{cell::OnceCell, time::Duration}; +use compact_str::CompactString; use jiff::{ Timestamp, civil::{Date, Time}, @@ -73,9 +74,9 @@ impl NexusKeysCfg { pub struct NexusHostCfg { pub base_url: url::Url, pub unix_path: Option<String>, - pub host_id: String, - pub user_id: String, - pub partner_id: String, + pub host_id: CompactString, + pub user_id: CompactString, + pub partner_id: CompactString, } impl NexusHostCfg { @@ -84,9 +85,9 @@ impl NexusHostCfg { Ok(Self { base_url: s.url("host_base_url").require()?, unix_path: s.path("UNIXPATH").opt()?, - host_id: s.str("host_id").require()?, - user_id: s.str("user_id").require()?, - partner_id: s.str("partner_id").require()?, + host_id: s.cstr("host_id").require()?, + user_id: s.cstr("user_id").require()?, + partner_id: s.cstr("partner_id").require()?, }) } @@ -233,6 +234,7 @@ pub struct NexusCfg { pub setup: OnceCell<NexusSetupConfig>, pub wire_cfg: Option<ApiCfg>, pub revenue_cfg: Option<ApiCfg>, + pub db_cfg: DbCfg, pub serve_cfg: Serve, } @@ -244,13 +246,14 @@ impl NexusCfg { account_type: s.parse("account type", "ACCOUNT_TYPE").require()?, wire_cfg: ApiCfg::parse(cfg.section("nexus-httpd-wire-gateway-api"))?, revenue_cfg: ApiCfg::parse(cfg.section("nexus-httpd-revenue-api"))?, - serve_cfg: Serve::parse(cfg.section("nexus-httpd"))?, + serve_cfg: Serve::parse(&cfg.section("nexus-httpd"))?, keys: OnceCell::new(), host: OnceCell::new(), fetch: OnceCell::new(), submit: OnceCell::new(), ebics: OnceCell::new(), setup: OnceCell::new(), + db_cfg: parse_db_cfg(&cfg)?, cfg, }) } diff --git a/libeufin-nexus/src/db.rs b/libeufin-nexus/src/db.rs @@ -14,14 +14,10 @@ TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> */ -use jiff::Timestamp; -use sqlx::{PgPool, types::Json}; -use taler_api::db::BindHelper; -use taler_common::config::Config; +use sqlx::PgPool; +use taler_api::config::DbCfg; use tokio::sync::watch::Sender; -use crate::{TaskStatus, config::parse_db_cfg}; - pub mod exchange; pub mod initiated; pub mod list; @@ -32,17 +28,15 @@ const SCHEMA: &str = "libeufin_nexus"; pub const UNSETTLED: &str = "status NOT IN ('success', 'permanent_failure', 'late_failure')"; pub const PENDING: &str = "status IN ('unsubmitted', 'pending')"; -pub async fn pool(cfg: &Config) -> anyhow::Result<PgPool> { - let db = parse_db_cfg(cfg)?; - let pool = taler_common::db::pool(db.cfg, SCHEMA).await?; +pub async fn pool(cfg: &DbCfg) -> anyhow::Result<PgPool> { + let pool = taler_common::db::pool(cfg.cfg.clone(), SCHEMA).await?; Ok(pool) } -pub async fn dbinit(cfg: &Config, reset: bool) -> anyhow::Result<PgPool> { - let db_cfg = parse_db_cfg(cfg)?; - let pool = taler_common::db::pool(db_cfg.cfg, SCHEMA).await?; +pub async fn dbinit(cfg: &DbCfg, reset: bool) -> anyhow::Result<PgPool> { + let pool = taler_common::db::pool(cfg.cfg.clone(), SCHEMA).await?; let mut db = pool.acquire().await?; - taler_common::db::dbinit(&mut db, db_cfg.sql_dir.as_ref(), "libeufin-nexus", reset).await?; + taler_common::db::dbinit(&mut db, cfg.sql_dir.as_ref(), "libeufin-nexus", reset).await?; Ok(pool) } @@ -65,33 +59,9 @@ pub async fn notification_listener( ) } -/** Get current value for [key] */ -pub async fn get_task_status(db: &PgPool, key: &str) -> sqlx::Result<Option<TaskStatus>> { - sqlx::query_scalar::<_, Json<TaskStatus>>("SELECT value FROM kv WHERE key=$1") - .bind(key) - .fetch_optional(db) - .await - .map(|it| it.map(|it| it.0)) -} - -/** Update a TaskStatus timestamp */ -pub async fn update_task_status( - db: &PgPool, - key: &str, - timestamp: &Timestamp, - success: bool, -) -> sqlx::Result<()> { - sqlx::query(if success { - "INSERT INTO kv (key, value) VALUES ($1, jsonb_build_object('last_successfull', $2, 'last_trial', $3)) ON CONFLICT (key) DO UPDATE SET value=EXCLUDED.value" - } else { - "INSERT INTO kv (key, value) VALUES ($1, jsonb_build_object('last_trial', $2)) ON CONFLICT (key) DO UPDATE SET value=jsonb_set(EXCLUDED.value, '{last_trial}'::text[], to_jsonb($3))" - }).bind(key).bind_timestamp(timestamp).bind_timestamp(timestamp).execute(db).await?; - Ok(()) -} - #[cfg(test)] pub mod test { - use libeufin_ebics::db::{ebics_first, ebics_register, ebics_remove}; + use libeufin_ebics::db::test::ebics_routine; use sqlx::{PgPool, Postgres, Row, pool::PoolConnection, postgres::PgRow}; use taler_api::db::TypeHelper; use taler_common::db::IncomingType; @@ -205,17 +175,8 @@ pub mod test { } #[tokio::test] - pub async fn ebics_pending() { + pub async fn ebics() { let (_, db) = db_setup().await; - let ids = ["first", "second", "third"]; - - for id in ids { - ebics_register(&db, id).await.unwrap(); - } - for id in ids { - assert_eq!(Some(id), ebics_first(&db).await.unwrap().as_deref()); - ebics_remove(&db, id).await.unwrap(); - } - assert_eq!(ebics_first(&db).await.unwrap(), None); + ebics_routine(&db).await; } } diff --git a/libeufin-nexus/src/fetch.rs b/libeufin-nexus/src/fetch.rs @@ -27,6 +27,7 @@ use anyhow::{anyhow, bail}; use compact_str::CompactStringExt; use jiff::{Timestamp, Zoned, tz::TimeZone}; use libeufin_ebics::{ + db::{get_task_status, update_task_status}, ebics::{ EbicsClient, EbicsErrKind, ebics_code::EbicsReturnCode, @@ -55,7 +56,6 @@ use crate::{ CHECKPOINT_KEY, FETCH_TASK_KEY, config::{AccountType, NexusCfg, NexusIngestCfg}, db::{ - get_task_status, initiated::{ batch_status_update, order_failure, order_step, order_success, tx_status_update, unsettled_tx_in_batch, @@ -65,7 +65,6 @@ use crate::{ OutgoingRegistrationResult, register_in, register_in_malformed, register_in_qr_bill, register_in_talerable, register_out_tx, }, - update_task_status, }, model::SubmissionState, rand_ebics_id, @@ -115,7 +114,7 @@ pub async fn ebics_fetch( } OrderDoc::status => { let msg_status = parse_pain002(&xml)?; - debug!(target: "ebics-fetch", "{msg_status}"); + debug!(target: "fetch", "{msg_status}"); if let Some(code) = msg_status.status { let msg = msg_status.msg(); batch_status_update( @@ -233,7 +232,7 @@ pub async fn ebics_fetch( match bank { EbicsReturnCode::EBICS_NO_DOWNLOAD_DATA_AVAILABLE => continue, EbicsReturnCode::EBICS_AUTHORISATION_ORDER_IDENTIFIER_FAILED => { - error!(target: "ebics-fetch", "{e}"); + error!(target: "fetch", "{e}"); success = false; continue; } @@ -268,8 +267,6 @@ pub async fn ebics_fetch( info!(target: "fetch", "Running with a frequency of {}", fetch_cfg.frequency_raw); } - // TODO loop - let mut last_fetch = Timestamp::UNIX_EPOCH; loop { let now = Timestamp::now(); diff --git a/libeufin-nexus/src/lib.rs b/libeufin-nexus/src/lib.rs @@ -17,13 +17,16 @@ * <http://www.gnu.org/licenses/> */ +#![allow(clippy::too_many_arguments)] + use std::{fmt::Write, sync::Arc, time::Duration}; use anyhow::{anyhow, bail}; use compact_str::{CompactString, CompactStringExt, ToCompactString}; use jiff::{Timestamp, civil::Date}; use libeufin_ebics::{ - cli::EbicsLogs, + cli::{EbicsArgs, EbicsLogs}, + db::update_task_status, ebics::{ EbicsClient, EbicsCtx, EbicsErrKind, EbicsError, EbicsErrorHelper as _, administrative::{AccountInfo, HKD, OrderInfo}, @@ -33,7 +36,6 @@ use libeufin_ebics::{ iso20022::pain001::{Pain001Msg, Pain001Tx, create_pain001}, keys::{BankKeys, ClientKeys}, }; -use serde::{Deserialize, Deserializer, Serialize, Serializer}; use sqlx::PgPool; use taler_api::api::TalerRouter; use taler_build::long_version; @@ -48,13 +50,13 @@ use tracing::{debug, error, info, warn}; use crate::{ api::NexusApi, - config::{NexusCfg, NexusEbicsConfig}, + config::{NexusCfg, NexusEbicsConfig, parse_db_cfg}, db::{ dbinit, initiated::{ batch_initiated, batch_sub_failure, batch_sub_success, initiate, initiated_submittable, }, - pool, update_task_status, + pool, }, fetch::ebics_fetch, list::ListCmd, @@ -83,16 +85,6 @@ const FETCH_TASK_KEY: &str = "fetch_task"; pub const CONFIG_SOURCE: ConfigSource = ConfigSource::new("libeufin", "libeufin-nexus", "libeufin-nexus"); -#[derive(clap::Parser, Debug, Clone)] -pub struct EbicsArgs { - #[command(flatten)] - logs: EbicsLogs, - - /// Execute once and return, ignoring the 'FREQUENCY' configuration value - #[arg(long)] - transient: bool, -} - #[derive(clap::Subcommand, Debug)] pub enum Cmd { /// Initialize libeufin-nexus database @@ -227,7 +219,7 @@ pub async fn ebics_submit( -> Result<CompactString, EbicsError> { let ctx = EbicsCtx::new(order); let xml = batch_pain001(batch, ebics_cfg, instant).ctx(&ctx)?; - ebics.upload(client, bank, order, &xml).await + ebics.upload(client, bank, order, xml.as_bytes()).await }; let submit_all = async || -> anyhow::Result<()> { @@ -454,6 +446,7 @@ pub async fn ebics_setup( pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { match cmd { Cmd::Dbinit { reset } => { + let cfg = parse_db_cfg(&cfg)?; dbinit(&cfg, reset).await?; } Cmd::EbicsSetup { @@ -462,8 +455,8 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { auto_accept_keys, generate_registration_pdf, } => { - let pool = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; let ebics = EbicsClient::new(cfg.host()?.ebics(), ebics_logs)?; ebics_setup( &ebics, @@ -481,8 +474,8 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { checkpoint, ebics: EbicsArgs { logs, transient }, } => { - let pool = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; let ebics = EbicsClient::new(cfg.host()?.ebics(), logs)?; let (client, bank) = cfg.expect_full_keys()?; ebics_fetch( @@ -502,8 +495,8 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { Cmd::EbicsSubmit { ebics: EbicsArgs { logs, transient }, } => { - let pool = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; let ebics = EbicsClient::new(cfg.host()?.ebics(), logs)?; let (client, bank) = cfg.expect_full_keys()?; ebics_submit(&ebics, &cfg, &client, &bank, &pool, transient).await? @@ -514,8 +507,8 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { end_to_end_id, payto, } => { - let pool = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; let subject = payto .subject @@ -567,8 +560,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { std::process::exit(1); } } else { - let pool = pool(&cfg.cfg).await?; - + let pool = pool(&cfg.db_cfg).await?; let payto = &cfg.ebics()?.account; let api = Arc::new(NexusApi::start(pool, payto.clone(), cfg.currency).await); let mut server = Router::new(); @@ -584,13 +576,13 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { } } Cmd::Manual(cmd) => { - let pool = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; cmd.run(&pool, &cfg).await?; } Cmd::List(cmd) => { - let pool = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let pool = pool(&cfg.db_cfg).await?; cmd.run(&pool, &cfg.currency).await?; } Cmd::Config(cmd) => cmd.run(&cfg)?, @@ -598,22 +590,3 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { } Ok(()) } - -#[derive(Debug, Serialize, Deserialize, Clone, Default)] -pub struct TaskStatus { - #[serde(serialize_with = "ser_micros", deserialize_with = "de_micros", default)] - pub last_successfull: Option<Timestamp>, - #[serde(serialize_with = "ser_micros", deserialize_with = "de_micros", default)] - pub last_trial: Option<Timestamp>, -} - -fn ser_micros<S: Serializer>(key: &Option<Timestamp>, serializer: S) -> Result<S::Ok, S::Error> { - key.map(|it| it.as_microsecond()).serialize(serializer) -} - -fn de_micros<'de, D: Deserializer<'de>>(deserializer: D) -> Result<Option<Timestamp>, D::Error> { - Option::<i64>::deserialize(deserializer)? - .map(Timestamp::from_microsecond) - .transpose() - .map_err(|e| serde::de::Error::custom(e.to_string())) -} diff --git a/libeufin-nexus/src/testing.rs b/libeufin-nexus/src/testing.rs @@ -135,8 +135,8 @@ impl TestingCmd { subject, payto, } => { - let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let db = pool(&cfg.db_cfg).await?; let subject = payto .subject .as_ref() @@ -170,8 +170,8 @@ impl TestingCmd { .await?; } TestingCmd::List(list_cmd) => { - let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let db = pool(&cfg.db_cfg).await?; list_cmd.run(&db, &cfg.currency).await?; } TestingCmd::EbicsBtd { @@ -187,8 +187,8 @@ impl TestingCmd { peek, dry_run, } => { - let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let db = pool(&cfg.db_cfg).await?; let order = Order::from_parts( &ty, Some(BTF { @@ -227,8 +227,8 @@ impl TestingCmd { .await?; } TestingCmd::TxCheck { logs } => { - let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let db = pool(&cfg.db_cfg).await?; let (client, bank) = cfg.expect_full_keys()?; let ebics = EbicsClient::new(cfg.host()?.ebics(), logs)?; let dialect = cfg.ebics()?.dialect.standard(); @@ -244,8 +244,8 @@ impl TestingCmd { println!("{res:?}") } TestingCmd::Wss { logs } => { - let db = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; + let db = pool(&cfg.db_cfg).await?; let (client, bank) = cfg.expect_full_keys()?; let ebics = EbicsClient::new(cfg.host()?.ebics(), logs)?; let (sender, mut receiver) = tokio::sync::mpsc::channel(10); diff --git a/testbench/Cargo.toml b/testbench/Cargo.toml @@ -17,7 +17,7 @@ taler-common.workspace = true libeufin-nexus = { path = "../libeufin-nexus"} libeufin-ebics = { path = "../libeufin-ebics"} reedline = "0.48" -shlex = "2.0" +shlex.workspace = true owo-colors = "4.3" tracing-subscriber = "0.3"