libeufin

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

commit 84704f78195558ee2a5d2ff7b125cb49a0ff435d
parent f5a57a3f1eb4e92b5d1f80eb3d68ccf2fd39e7a7
Author: Antoine A <>
Date:   Fri, 24 Apr 2026 10:37:58 +0200

nexus: rust taler API

Diffstat:
MCargo.lock | 95++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------------
Asrc/api.rs | 457+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/db.rs | 234+++++++------------------------------------------------------------------------
Asrc/db/exchange.rs | 325+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/db/payment.rs | 80+++++++++++++++++++++++++++++++++++++++----------------------------------------
Asrc/db/transfer.rs | 83+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/lib.rs | 1+
Msrc/model.rs | 14+++++++++++++-
8 files changed, 997 insertions(+), 292 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -150,9 +150,9 @@ checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" [[package]] name = "aws-lc-rs" -version = "1.16.1" +version = "1.16.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "94bffc006df10ac2a68c83692d734a465f8ee6c5b384d8545a636f81d858f4bf" +checksum = "a054912289d18629dc78375ba2c3726a3afe3ff71b4edba9dedfca0e3446d1fc" dependencies = [ "aws-lc-sys", "untrusted 0.7.1", @@ -161,9 +161,9 @@ dependencies = [ [[package]] name = "aws-lc-sys" -version = "0.38.0" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4321e568ed89bb5a7d291a7f37997c2c0df89809d7b6d12062c81ddb54aa782e" +checksum = "1fa7e52a4c5c547c741610a2c6f123f3881e409b714cd27e6798ef020c514f0a" dependencies = [ "cc", "cmake", @@ -282,9 +282,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.2.57" +version = "1.2.58" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7a0dd1ca384932ff3641c8718a02769f1698e7563dc6974ffd03346116310423" +checksum = "e1e928d4b69e3077709075a938a05ffbedfa53a84c8f766efbf8220bb1ff60e1" dependencies = [ "find-msvc-tools", "jobserver", @@ -375,9 +375,9 @@ checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" [[package]] name = "cmake" -version = "0.1.57" +version = "0.1.58" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "75443c44cd6b379beb8c5b45d85d0773baf31cce901fe7bb252f4eff3008ef7d" +checksum = "c0f78a02292a74a88ac736019ab962ece0bc380e3f977bf72e376c5d78ff0678" dependencies = [ "cc", ] @@ -1337,9 +1337,9 @@ checksum = "d98f6fed1fde3f8c21bc40a1abb88dd75e67924f9cffc3ef95607bad8017f8e2" [[package]] name = "iri-string" -version = "0.7.10" +version = "0.7.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c91338f0783edbd6195decb37bae672fd3b165faffb89bf7b9e6942f8b1a731a" +checksum = "d8e7418f59cc01c88316161279a7f665217ae316b388e58a0d10e29f54f1e5eb" dependencies = [ "memchr", "serde", @@ -1362,9 +1362,9 @@ dependencies = [ [[package]] name = "itoa" -version = "1.0.17" +version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "92ecc6618181def0457392ccd0ee51198e065e016d1d527a7ac1b6dc7c1f09d2" +checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" [[package]] name = "jiff" @@ -1416,7 +1416,7 @@ dependencies = [ "cesu8", "cfg-if", "combine", - "jni-sys", + "jni-sys 0.3.1", "log", "thiserror 1.0.69", "walkdir", @@ -1425,9 +1425,31 @@ dependencies = [ [[package]] name = "jni-sys" -version = "0.3.0" +version = "0.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8eaf4bc02d17cbdd7ff4c7438cafcdf7fb9a4613313ad11b4f8fefe7d3fa0130" +checksum = "41a652e1f9b6e0275df1f15b32661cf0d4b78d4d87ddec5e0c3c20f097433258" +dependencies = [ + "jni-sys 0.4.1", +] + +[[package]] +name = "jni-sys" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c6377a88cb3910bee9b0fa88d4f42e1d2da8e79915598f65fb0c7ee14c878af2" +dependencies = [ + "jni-sys-macros", +] + +[[package]] +name = "jni-sys-macros" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "38c0b942f458fe50cdac086d2f946512305e5631e720728f2a61aabcd47a6264" +dependencies = [ + "quote", + "syn", +] [[package]] name = "jobserver" @@ -1530,9 +1552,9 @@ checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981" [[package]] name = "libredox" -version = "0.1.14" +version = "0.1.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1744e39d1d6a9948f4f388969627434e31128196de472883b39f148769bfe30a" +checksum = "7ddbf48fd451246b1f8c2610bd3b4ac0cc6e149d89832867093ab69a17194f08" dependencies = [ "bitflags", "libc", @@ -1703,9 +1725,9 @@ dependencies = [ [[package]] name = "num-conv" -version = "0.2.0" +version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cf97ec579c3c42f953ef76dbf8d55ac91fb219dde70e49aa4a6b7d74e9919050" +checksum = "c6673768db2d862beb9b39a78fdcb1a69439615d5794a1be50caa9bc92c81967" [[package]] name = "num-integer" @@ -2339,9 +2361,9 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.9" +version = "0.103.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d7df23109aa6c1567d1c575b9952556388da57401e4ace1d15f79eedad0d8f53" +checksum = "df33b2b81ac578cabaf06b89b0631153a3f416b0a886e8a7a1707fb51abbd1ef" dependencies = [ "aws-lc-rs", "ring", @@ -2582,9 +2604,9 @@ dependencies = [ [[package]] name = "simd-adler32" -version = "0.3.8" +version = "0.3.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e320a6c5ad31d271ad523dcf3ad13e2767ad8b1cb8f047f75a8aeaf8da139da2" +checksum = "703d5c7ef118737c72f1af64ad2f6f8c5e1921f818cdcb97b8fe6fc69bf66214" [[package]] name = "slab" @@ -2959,7 +2981,7 @@ dependencies = [ [[package]] name = "taler-api" -version = "1.4.0" +version = "1.5.0" dependencies = [ "aws-lc-rs", "axum", @@ -2969,6 +2991,7 @@ dependencies = [ "http-body-util", "jiff", "listenfd", + "regex", "serde", "serde_json", "serde_path_to_error", @@ -2983,11 +3006,11 @@ dependencies = [ [[package]] name = "taler-build" -version = "1.4.0" +version = "1.5.0" [[package]] name = "taler-common" -version = "1.4.0" +version = "1.5.0" dependencies = [ "anyhow", "aws-lc-rs", @@ -3014,11 +3037,13 @@ dependencies = [ [[package]] name = "taler-test-utils" -version = "1.4.0" +version = "1.5.0" dependencies = [ + "aws-lc-rs", "axum", "flate2", "http-body-util", + "jiff", "serde", "serde_json", "serde_urlencoded", @@ -3362,9 +3387,9 @@ checksum = "7df058c713841ad818f1dc5d3fd88063241cc61f49f5fbea4b951e8cf5a8d71d" [[package]] name = "unicode-segmentation" -version = "1.12.0" +version = "1.13.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f6ccf251212114b54433ec949fd6a7841275f9ada20dddd2f29e9ceea4501493" +checksum = "9629274872b2bfaf8d66f5f15725007f635594914870f65218920345aa11aa8c" [[package]] name = "unicode-width" @@ -3417,9 +3442,9 @@ checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" [[package]] name = "uuid" -version = "1.22.0" +version = "1.23.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a68d3c8f01c0cfa54a75291d83601161799e4a89a39e0929f4b0354d88757a37" +checksum = "5ac8b6f42ead25368cf5b098aeb3dc8a1a2c05a3eee8a9a1a68c640edbfc79d9" dependencies = [ "getrandom 0.4.2", "js-sys", @@ -4197,18 +4222,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.42" +version = "0.8.47" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f2578b716f8a7a858b7f02d5bd870c14bf4ddbbcf3a4c05414ba6503640505e3" +checksum = "efbb2a062be311f2ba113ce66f697a4dc589f85e78a4aea276200804cea0ed87" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.42" +version = "0.8.47" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7e6cc098ea4d3bd6246687de65af3f920c430e236bee1e3bf2e441463f08a02f" +checksum = "0e8bc7269b54418e7aeeef514aa68f8690b8c0489a06b0136e5f57c4c5ccab89" dependencies = [ "proc-macro2", "quote", diff --git a/src/api.rs b/src/api.rs @@ -0,0 +1,457 @@ +/* +* 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 jiff::Timestamp; +use taler_api::{ + api::{TalerApi, revenue::Revenue, transfer::WireTransferGateway, wire::WireGateway}, + error::{ApiResult, failure, failure_code}, + subject::{IncomingSubject, fmt_in_subject, subject_fmt_qr_bill}, +}; +use taler_common::{ + api_common::{SafeU64, safe_u64}, + api_params::{History, Page}, + api_revenue::RevenueIncomingHistory, + api_transfer::{ + RegistrationRequest, RegistrationResponse, SubjectFormat, TransferSubject, Unregistration, + }, + api_wire::{ + AddIncomingRequest, AddIncomingResponse, AddKycauthRequest, AddKycauthResponse, + IncomingHistory, OutgoingHistory, TransferList, TransferRequest, TransferResponse, + TransferState, TransferStatus, + }, + db::IncomingType, + error_code::ErrorCode, + types::{ + amount::Currency, + payto::{FullIbanPayto, PaytoURI}, + timestamp::TalerTimestamp, + }, +}; +use tokio::sync::watch::Sender; + +use crate::{ + db::{ + self, + exchange::{ + TransferResult, incoming_history, outgoing_history, revenue_history, transfer, + transfer_by_id, transfer_page, + }, + payment::{IncomingRegistrationResult, register_in_talerable}, + transfer::{RegistrationResult, transfer_register, transfer_unregister}, + }, + model::{IncomingId, IncomingPayment}, + rand_ebics_id, +}; + +pub struct NexusApi { + pub pool: sqlx::PgPool, + pub currency: Currency, + pub payto: PaytoURI, + pub in_channel: Sender<i64>, + pub taler_in_channel: Sender<i64>, + pub taler_out_channel: Sender<i64>, +} + +impl NexusApi { + pub async fn start(pool: sqlx::PgPool, payto: PaytoURI, currency: Currency) -> Self { + let in_channel = Sender::new(0); + let taler_in_channel = Sender::new(0); + let taler_out_channel = Sender::new(0); + let tmp = Self { + pool: pool.clone(), + payto, + currency, + in_channel: in_channel.clone(), + taler_in_channel: taler_in_channel.clone(), + taler_out_channel: taler_out_channel.clone(), + }; + tokio::spawn(db::notification_listener( + pool, + in_channel, + taler_in_channel, + taler_out_channel, + )); + tmp + } +} + +impl TalerApi for NexusApi { + fn currency(&self) -> &str { + self.currency.as_ref() + } + + fn implementation(&self) -> &'static str { + "urn:net:taler:specs:libeufin-nexus:taler-rust" + } +} + +impl WireGateway for NexusApi { + async fn transfer(&self, req: TransferRequest) -> ApiResult<TransferResponse> { + FullIbanPayto::try_from(&req.credit_account)?; + let result = transfer(&self.pool, &req, &rand_ebics_id(), &Timestamp::now()).await?; + match result { + TransferResult::Success { id, timestamp } => Ok(TransferResponse { + timestamp: timestamp.into(), + row_id: SafeU64::try_from(id).unwrap(), + }), + TransferResult::RequestUidReuse => { + Err(failure_code(ErrorCode::BANK_TRANSFER_REQUEST_UID_REUSED)) + } + TransferResult::WtidReuse => Err(failure_code(ErrorCode::BANK_TRANSFER_WTID_REUSED)), + } + } + + async fn transfer_page( + &self, + page: Page, + status: Option<TransferState>, + ) -> ApiResult<TransferList> { + Ok(TransferList { + transfers: transfer_page(&self.pool, &self.currency, &page, &status).await?, + debit_account: self.payto.clone(), + }) + } + + async fn transfer_by_id(&self, id: u64) -> ApiResult<Option<TransferStatus>> { + Ok(transfer_by_id(&self.pool, &self.currency, id).await?) + } + + async fn outgoing_history(&self, params: History) -> ApiResult<OutgoingHistory> { + Ok(OutgoingHistory { + outgoing_transactions: outgoing_history(&self.pool, &self.currency, &params, || { + self.taler_out_channel.subscribe() + }) + .await?, + debit_account: self.payto.clone(), + }) + } + + async fn incoming_history(&self, params: History) -> ApiResult<IncomingHistory> { + Ok(IncomingHistory { + incoming_transactions: incoming_history(&self.pool, &self.currency, &params, || { + self.taler_in_channel.subscribe() + }) + .await?, + credit_account: self.payto.clone(), + }) + } + + async fn add_incoming_reserve( + &self, + req: AddIncomingRequest, + ) -> ApiResult<AddIncomingResponse> { + FullIbanPayto::try_from(&req.debit_account)?; + let now = Timestamp::now(); + match register_in_talerable( + &self.pool, + &IncomingPayment { + id: IncomingId { + uetr: None, + tx_id: Some(rand_ebics_id()), + acct_svcr_ref: None, + }, + amount: req.amount, + credit_fee: None, + subject: Some(format!( + "Manual incoming {}", + fmt_in_subject(IncomingType::reserve, &req.reserve_pub) + )), + execution_time: now, + debtor: Some(req.debit_account), + }, + &IncomingSubject::Reserve(req.reserve_pub), + ) + .await? + { + IncomingRegistrationResult::Success(in_result) => Ok(AddIncomingResponse { + row_id: safe_u64(in_result.id), + timestamp: now.into(), + }), + IncomingRegistrationResult::ReservePubReuse => { + Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) + } + IncomingRegistrationResult::MappingReuse + | IncomingRegistrationResult::UnknownMapping => unreachable!("mapping not used"), + } + } + + async fn add_incoming_kyc(&self, req: AddKycauthRequest) -> ApiResult<AddKycauthResponse> { + FullIbanPayto::try_from(&req.debit_account)?; + let now = Timestamp::now(); + match register_in_talerable( + &self.pool, + &IncomingPayment { + id: IncomingId { + uetr: None, + tx_id: Some(rand_ebics_id()), + acct_svcr_ref: None, + }, + amount: req.amount, + credit_fee: None, + subject: Some(format!( + "Manual incoming {}", + fmt_in_subject(IncomingType::kyc, &req.account_pub) + )), + execution_time: now, + debtor: Some(req.debit_account), + }, + &IncomingSubject::Kyc(req.account_pub), + ) + .await? + { + IncomingRegistrationResult::Success(in_result) => Ok(AddIncomingResponse { + row_id: safe_u64(in_result.id), + timestamp: now.into(), + }), + IncomingRegistrationResult::ReservePubReuse => { + Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) + } + IncomingRegistrationResult::MappingReuse + | IncomingRegistrationResult::UnknownMapping => unreachable!("mapping not used"), + } + } + + fn support_account_check(&self) -> bool { + false + } +} + +impl Revenue for NexusApi { + async fn history(&self, params: History) -> ApiResult<RevenueIncomingHistory> { + Ok(RevenueIncomingHistory { + incoming_transactions: revenue_history(&self.pool, &self.currency, &params, || { + self.in_channel.subscribe() + }) + .await?, + credit_account: self.payto.clone(), + }) + } +} + +impl WireTransferGateway for NexusApi { + fn supported_formats(&self) -> &[SubjectFormat] { + &[SubjectFormat::SIMPLE] + } + + async fn registration(&self, req: RegistrationRequest) -> ApiResult<RegistrationResponse> { + let reference_number = subject_fmt_qr_bill(req.authorization_pub.as_ref()); + match transfer_register( + &self.pool, + req.r#type.into(), + &req.account_pub, + &req.authorization_pub, + &req.authorization_sig, + req.recurrent, + &reference_number, + &Timestamp::now(), + ) + .await? + { + RegistrationResult::Success => ApiResult::Ok(RegistrationResponse { + subjects: vec![ + TransferSubject::QrBill { + credit_amount: req.credit_amount.clone(), + qr_reference_number: reference_number, + }, + TransferSubject::Simple { + credit_amount: req.credit_amount, + subject: if req.authorization_pub == req.account_pub && !req.recurrent { + fmt_in_subject(req.r#type.into(), &req.account_pub) + } else { + fmt_in_subject(IncomingType::map, &req.authorization_pub) + }, + }, + ], + expiration: TalerTimestamp::Never, + }), + RegistrationResult::ReservePubReuse => { + ApiResult::Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) + } + RegistrationResult::SubjectReuse => { + ApiResult::Err(failure_code(ErrorCode::BANK_DERIVATION_REUSE)) + } + } + } + + async fn unregistration(&self, req: Unregistration) -> ApiResult<()> { + if !transfer_unregister(&self.pool, &req.authorization_pub, &Timestamp::now()).await? { + Err(failure( + ErrorCode::BANK_TRANSACTION_NOT_FOUND, + format!("Prepared transfer '{}' not found", req.authorization_pub), + )) + } else { + Ok(()) + } + } +} + +#[cfg(test)] +mod test { + use std::{ + str::FromStr as _, + sync::{Arc, LazyLock}, + }; + + use sqlx::PgPool; + use taler_api::{ + api::TalerRouter as _, + auth::AuthMethod, + subject::{OutgoingSubject, fmt_in_subject}, + }; + use taler_common::{ + api_revenue::RevenueConfig, + api_transfer::WireTransferConfig, + api_wire::{OutgoingHistory, TransferState, WireConfig}, + db::IncomingType, + types::{ + amount::Currency, + payto::{PaytoURI, payto}, + }, + }; + use taler_test_utils::{ + Router, + db::db_test_setup, + routine::{ + admin_add_incoming_routine, registration_routine, revenue_routine, routine_pagination, + transfer_routine, + }, + server::TestServer as _, + }; + + use crate::{ + CONFIG_SOURCE, + api::NexusApi, + config::{AccountType, NexusIngestConfig}, + db::{ + payment::register_out_tx, + test::{CURRENCY, check_in, gen_in_pay, gen_out_pay}, + }, + worker::register_incoming, + }; + + static ACCOUNT: LazyLock<PaytoURI> = + LazyLock::new(|| payto("payto://iban/CH4189144589712575493?receiver-name=Test")); + + async fn setup() -> (Router, PgPool) { + let (_, pool) = db_test_setup(CONFIG_SOURCE).await; + let api = Arc::new( + NexusApi::start( + pool.clone(), + ACCOUNT.clone(), + Currency::from_str("TEST").unwrap(), + ) + .await, + ); + let server = Router::new() + .wire_gateway(api.clone(), AuthMethod::None) + .wire_transfer_gateway(api.clone()) + .revenue(api, AuthMethod::None) + .finalize(); + + (server, pool) + } + + #[tokio::test] + async fn config() { + let (server, _) = setup().await; + server + .get("/taler-wire-gateway/config") + .await + .assert_ok_json::<WireConfig>(); + server + .get("/taler-wire-transfer-gateway/config") + .await + .assert_ok_json::<WireTransferConfig>(); + server + .get("/taler-revenue/config") + .await + .assert_ok_json::<RevenueConfig>(); + } + + #[tokio::test] + async fn transfer() { + let (server, _) = setup().await; + transfer_routine(&server, TransferState::pending, &ACCOUNT).await; + // TODO + /*db.initiated.batchSubmissionSuccess(1, Instant.now(), "ORDER1") + db.initiated.batchSubmissionFailure(2, Instant.now(), "Failure") + db.initiated.batchSubmissionFailure(3, Instant.now(), "Failure") + client.getA("/taler-wire-gateway/transfers?status=transient_failure").assertOkJson<TransferList> { + assertEquals(2, it.transfers.size) + } + client.getA("/taler-wire-gateway/transfers?status=pending").assertOkJson<TransferList> { + assertEquals(4, it.transfers.size) + }*/ + } + + #[tokio::test] + async fn outgoing_history() { + let (server, pool) = setup().await; + routine_pagination::<OutgoingHistory>( + &server, + "/taler-wire-gateway/history/outgoing", + async |_| { + register_out_tx( + &pool, + &gen_out_pay("subject"), + Some(&OutgoingSubject::rand()), + ) + .await + .unwrap(); + }, + ) + .await; + } + + #[tokio::test] + async fn admin_add_incoming() { + let (server, _) = setup().await; + admin_add_incoming_routine(&server, &ACCOUNT, true).await; + } + + #[tokio::test] + async fn revenue() { + let (server, _) = setup().await; + revenue_routine(&server, &ACCOUNT, true).await; + } + + #[tokio::test] + async fn registration() { + let (server, pool) = setup().await; + registration_routine( + &server, + &ACCOUNT, + || check_in(&pool), + |account_pub| { + let account_pub = account_pub.clone(); + let pool = &pool; + async move { + register_incoming( + pool, + &NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY), + &gen_in_pay(fmt_in_subject(IncomingType::map, &account_pub)), + ) + .await + .unwrap() + } + }, + ) + .await; + } +} diff --git a/src/db.rs b/src/db.rs @@ -15,22 +15,16 @@ */ use compact_str::CompactString; -use jiff::Timestamp; -use sqlx::{PgPool, QueryBuilder, Row, postgres::PgRow}; -use taler_api::db::{BindHelper, IncomingType, TypeHelper, history}; -use taler_common::{ - api_common::{EddsaPublicKey, EddsaSignature}, - api_params::History, - api_wire::{IncomingBankTransaction, OutgoingBankTransaction}, - config::Config, - types::amount::{Amount, Currency}, -}; -use tokio::sync::watch::{Receiver, Sender}; +use sqlx::{PgPool, Row, postgres::PgRow}; +use taler_common::config::Config; +use tokio::sync::watch::Sender; use crate::config::parse_db_cfg; +pub mod exchange; pub mod initiated; pub mod payment; +pub mod transfer; const SCHEMA: &str = "libeufin_nexus"; pub const UNSETTLED: &str = "status NOT IN ('success', 'permanent_failure', 'late_failure')"; @@ -69,188 +63,6 @@ pub async fn notification_listener( ) } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum RegistrationResult { - Success, - ReservePubReuse, - SubjectReuse, -} - -pub async fn transfer_register( - db: &PgPool, - ty: IncomingType, - account_pub: &EddsaPublicKey, - auth_pub: &EddsaPublicKey, - auth_sig: &EddsaSignature, - recurrent: bool, - reference_number: &str, - timestamp: &Timestamp, -) -> sqlx::Result<RegistrationResult> { - sqlx::query( - " - SELECT - out_subject_reuse, - out_reserve_pub_reuse - FROM register_prepared_transfers ( - $1::taler_incoming_type,$2,$3,$4,$5,$6,$7 - ) - ", - ) - .bind(ty) - .bind(account_pub) - .bind(auth_pub) - .bind(auth_sig) - .bind(recurrent) - .bind(reference_number) - .bind_timestamp(timestamp) - .try_map(|r: PgRow| { - Ok(if r.try_get_flag("out_subject_reuse")? { - RegistrationResult::SubjectReuse - } else if r.try_get_flag("out_reserve_pub_reuse")? { - RegistrationResult::ReservePubReuse - } else { - RegistrationResult::Success - }) - }) - .fetch_one(db) - .await -} - -pub async fn transfer_unregister( - db: &PgPool, - auth_pub: &EddsaPublicKey, - timestamp: &Timestamp, -) -> sqlx::Result<bool> { - sqlx::query("SELECT out_found FROM delete_prepared_transfers($1,$2)") - .bind(auth_pub) - .bind_timestamp(timestamp) - .try_map(|r: PgRow| r.try_get_flag(0)) - .fetch_one(db) - .await -} - -pub async fn outgoing_history( - db: &PgPool, - currency: &Currency, - params: &History, - listen: impl FnOnce() -> Receiver<i64>, -) -> sqlx::Result<Vec<OutgoingBankTransaction>> { - history( - db, - "outgoing_transaction_id", - params, - listen, - || { - QueryBuilder::new( - " - SELECT - outgoing_transaction_id - ,execution_time - ,(amount).val AS amount_val - ,(amount).frac AS amount_frac - ,(debit_fee).val AS debit_fee_val - ,(debit_fee).frac AS debit_fee_frac - ,credit_payto - ,wtid - ,exchange_base_url - ,metadata - FROM talerable_outgoing_transactions - JOIN outgoing_transactions USING(outgoing_transaction_id) - WHERE - ", - ) - }, - |r: PgRow| { - let tmp: Amount = r.try_get_amount("debit_fee", currency)?; - - Ok(OutgoingBankTransaction { - row_id: r.try_get_safeu64("outgoing_transaction_id")?, - amount: r.try_get_amount("amount", currency)?, - debit_fee: if tmp == Amount::zero(currency) { - None - } else { - Some(tmp) - }, - credit_account: r.try_get_payto("credit_payto")?, - date: r.try_get_timestamp("execution_time")?.into(), - exchange_base_url: r.try_get_url("exchange_base_url")?, - wtid: r.try_get("wtid")?, - }) - }, - ) - .await -} - -pub async fn incoming_history( - db: &PgPool, - currency: &Currency, - params: &History, - listen: impl FnOnce() -> Receiver<i64>, -) -> sqlx::Result<Vec<IncomingBankTransaction>> { - history( - db, - "incoming_transaction_id", - params, - listen, - || { - QueryBuilder::new( - " - SELECT - incoming_transaction_id - ,execution_time - ,(amount).val AS amount_val - ,(amount).frac AS amount_frac - ,(credit_fee).val AS credit_fee_val - ,(credit_fee).frac AS credit_fee_frac - ,debit_payto - ,type - ,metadata - ,authorization_pub - ,authorization_sig - FROM talerable_incoming_transactions - JOIN incoming_transactions USING(incoming_transaction_id) - WHERE - ", - ) - }, - |r: PgRow| { - let tmp: Amount = r.try_get_amount("credit_fee", currency)?; - Ok(match r.try_get("type")? { - IncomingType::reserve => IncomingBankTransaction::Reserve { - row_id: r.try_get_safeu64("incoming_transaction_id")?, - amount: r.try_get_amount("amount", currency)?, - credit_fee: if tmp == Amount::zero(currency) { - None - } else { - Some(tmp) - }, - debit_account: r.try_get_payto("debit_payto")?, - date: r.try_get_timestamp("execution_time")?.into(), - reserve_pub: r.try_get("metadata")?, - authorization_pub: r.try_get("authorization_pub")?, - authorization_sig: r.try_get("authorization_sig")?, - }, - IncomingType::kyc => IncomingBankTransaction::Kyc { - row_id: r.try_get_safeu64("incoming_transaction_id")?, - amount: r.try_get_amount("amount", currency)?, - credit_fee: if tmp == Amount::zero(currency) { - None - } else { - Some(tmp) - }, - debit_account: r.try_get_payto("debit_payto")?, - date: r.try_get_timestamp("execution_time")?.into(), - account_pub: r.try_get("metadata")?, - authorization_pub: r.try_get("authorization_pub")?, - authorization_sig: r.try_get("authorization_sig")?, - }, - IncomingType::map => unimplemented!("MAP are never listed in the history"), - }) - }, - ) - .await -} - /** Register a pending transaction */ pub async fn ebics_register(db: &PgPool, id: &str) -> sqlx::Result<()> { sqlx::query( @@ -280,20 +92,21 @@ pub async fn ebics_first(db: &PgPool) -> sqlx::Result<Option<CompactString>> { } #[cfg(test)] -mod test { +pub mod test { use std::{str::FromStr, sync::LazyLock}; use compact_str::CompactString; use jiff::Timestamp; - use sqlx::{PgConnection, PgPool, Row, postgres::PgRow}; - use taler_api::db::{IncomingType, TypeHelper}; + use sqlx::{PgPool, Postgres, Row, pool::PoolConnection, postgres::PgRow}; + use taler_api::db::TypeHelper; use taler_common::{ - api_common::EddsaPublicKey, + db::IncomingType, types::{ amount::{Amount, Currency}, payto::IbanPayto, }, }; + use taler_test_utils::routine::Status; use crate::{ CONFIG_SOURCE, @@ -304,10 +117,8 @@ mod test { pub static CURRENCY: LazyLock<Currency> = LazyLock::new(|| "KUDOS".parse().unwrap()); - pub async fn setup() -> (PgConnection, PgPool) { - let pool = taler_test_utils::db::db_test_setup(CONFIG_SOURCE).await; - let conn = pool.acquire().await.unwrap().leak(); - (conn, pool) + pub async fn setup() -> (PoolConnection<Postgres>, PgPool) { + taler_test_utils::db::db_test_setup(CONFIG_SOURCE).await } /** Generates an outgoing payment, given its subject */ @@ -426,18 +237,8 @@ mod test { .unwrap(); } - #[derive(Debug, PartialEq, Eq)] - pub enum Status { - Simple, - Pending, - Bounced, - Incomplete, - Reserve(EddsaPublicKey), - Kyc(EddsaPublicKey), - } - - pub async fn check_in(db: &PgPool, state: &[Status]) { - let current = sqlx::query( + pub async fn check_in(db: &PgPool) -> Vec<Status> { + sqlx::query( " SELECT pending_recurrent_incoming_transactions.authorization_pub IS NOT NULL, initiated_outgoing_transaction_id IS NOT NULL, debit_payto IS NULL OR subject IS NULL, type::text, metadata FROM incoming_transactions @@ -467,8 +268,11 @@ mod test { }) .fetch_all(db) .await - .unwrap(); - assert_eq!(state, current); + .unwrap() + } + + pub async fn check_in_state(db: &PgPool, state: &[Status]) { + assert_eq!(state, check_in(db).await); } #[tokio::test] diff --git a/src/db/exchange.rs b/src/db/exchange.rs @@ -0,0 +1,325 @@ +/* + This file is part of TALER + Copyright (C) 2026 Taler Systems SA + + TALER is free software; you can redistribute it and/or modify it under the + terms of the GNU Affero General Public License as published by the Free Software + Foundation; either version 3, or (at your option) any later version. + + TALER is distributed in the hope that it will be useful, but WITHOUT ANY + WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR + A PARTICULAR PURPOSE. See the GNU Affero General Public License for more details. + + You should have received a copy of the GNU Affero General Public License along with + TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> +*/ + +use jiff::Timestamp; +use sqlx::{PgPool, QueryBuilder, Row as _, postgres::PgRow}; +use taler_api::{ + db::{BindHelper, TypeHelper as _, history, page}, + serialized, + subject::fmt_out_subject, +}; +use taler_common::{ + api_params::{History, Page}, + api_revenue::RevenueIncomingBankTransaction, + api_wire::{ + IncomingBankTransaction, OutgoingBankTransaction, TransferListStatus, TransferRequest, + TransferState, TransferStatus, + }, + db::IncomingType, + types::amount::Currency, +}; +use tokio::sync::watch::Receiver; + +use crate::model::SubmissionState; + +pub async fn outgoing_history( + db: &PgPool, + currency: &Currency, + params: &History, + listen: impl FnOnce() -> Receiver<i64>, +) -> sqlx::Result<Vec<OutgoingBankTransaction>> { + history( + db, + "outgoing_transaction_id", + params, + listen, + || { + QueryBuilder::new( + " + SELECT + outgoing_transaction_id + ,execution_time + ,amount + ,debit_fee + ,credit_payto + ,wtid + ,exchange_base_url + ,metadata + FROM talerable_outgoing_transactions + JOIN outgoing_transactions USING(outgoing_transaction_id) + WHERE + ", + ) + }, + |r: PgRow| { + Ok(OutgoingBankTransaction { + row_id: r.try_get_safeu64("outgoing_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + debit_fee: r + .try_get_opt_amount("debit_fee", currency)? + .filter(|it| it.is_zero()), + credit_account: r.try_get_payto("credit_payto")?, + date: r.try_get_timestamp("execution_time")?.into(), + exchange_base_url: r.try_get_url("exchange_base_url")?, + wtid: r.try_get("wtid")?, + metadata: r.try_get("metadata")?, + }) + }, + ) + .await +} + +pub async fn incoming_history( + db: &PgPool, + currency: &Currency, + params: &History, + listen: impl FnOnce() -> Receiver<i64>, +) -> sqlx::Result<Vec<IncomingBankTransaction>> { + history( + db, + "incoming_transaction_id", + params, + listen, + || { + QueryBuilder::new( + " + SELECT + incoming_transaction_id + ,execution_time + ,amount + ,credit_fee + ,debit_payto + ,type::text + ,metadata + ,authorization_pub + ,authorization_sig + FROM talerable_incoming_transactions + JOIN incoming_transactions USING(incoming_transaction_id) + WHERE + ", + ) + }, + |r: PgRow| { + let credit_fee = r + .try_get_opt_amount("credit_fee", currency)? + .filter(|it| it.is_zero()); + Ok(match r.try_get("type")? { + IncomingType::reserve => IncomingBankTransaction::Reserve { + row_id: r.try_get_safeu64("incoming_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + credit_fee, + debit_account: r.try_get_payto("debit_payto")?, + date: r.try_get_timestamp("execution_time")?.into(), + reserve_pub: r.try_get("metadata")?, + authorization_pub: r.try_get("authorization_pub")?, + authorization_sig: r.try_get("authorization_sig")?, + }, + IncomingType::kyc => IncomingBankTransaction::Kyc { + row_id: r.try_get_safeu64("incoming_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + credit_fee, + debit_account: r.try_get_payto("debit_payto")?, + date: r.try_get_timestamp("execution_time")?.into(), + account_pub: r.try_get("metadata")?, + authorization_pub: r.try_get("authorization_pub")?, + authorization_sig: r.try_get("authorization_sig")?, + }, + IncomingType::map => unimplemented!("MAP are never listed in the history"), + }) + }, + ) + .await +} + +pub async fn revenue_history( + db: &PgPool, + currency: &Currency, + params: &History, + listen: impl FnOnce() -> Receiver<i64>, +) -> sqlx::Result<Vec<RevenueIncomingBankTransaction>> { + history( + db, + "incoming_transaction_id", + params, + listen, + || { + QueryBuilder::new( + " + SELECT + incoming_transaction_id + ,execution_time + ,amount + ,credit_fee + ,debit_payto + ,subject + FROM incoming_transactions + WHERE debit_payto IS NOT NULL AND subject IS NOT NULL AND + ", + ) + }, + |r: PgRow| { + Ok(RevenueIncomingBankTransaction { + row_id: r.try_get_safeu64("incoming_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + credit_fee: r + .try_get_opt_amount("credit_fee", currency)? + .filter(|it| it.is_zero()), + debit_account: r.try_get_payto("debit_payto")?, + date: r.try_get_timestamp("execution_time")?.into(), + subject: r.try_get("subject")?, + }) + }, + ) + .await +} + +pub enum TransferResult { + Success { id: u64, timestamp: Timestamp }, + RequestUidReuse, + WtidReuse, +} + +pub async fn transfer( + db: &PgPool, + req: &TransferRequest, + e2e_id: &str, + timestamp: &Timestamp, +) -> sqlx::Result<TransferResult> { + let subject = fmt_out_subject(&req.wtid, &req.exchange_base_url, req.metadata.as_deref()); + serialized!( + sqlx::query( + " + SELECT + out_request_uid_reuse + ,out_wtid_reuse + ,out_tx_row_id + ,out_timestamp + FROM taler_transfer($1,$2,$3,$4,$5,$6,$7,$8,$9) + ", + ) + .bind(&req.request_uid) + .bind(&req.wtid) + .bind(&subject) + .bind(&req.amount) + .bind(req.exchange_base_url.as_str()) + .bind(&req.metadata) + .bind(req.credit_account.as_ref().as_str()) + .bind(e2e_id) + .bind_timestamp(timestamp) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_request_uid_reuse")? { + TransferResult::RequestUidReuse + } else if r.try_get_flag("out_wtid_reuse")? { + TransferResult::WtidReuse + } else { + TransferResult::Success { + id: r.try_get_u64("out_tx_row_id")?, + timestamp: r.try_get_timestamp("out_timestamp")?, + } + }) + }) + .fetch_one(db) + ) +} + +pub async fn transfer_by_id( + db: &PgPool, + currency: &Currency, + id: u64, +) -> sqlx::Result<Option<TransferStatus>> { + serialized!( + sqlx::query( + " + SELECT + wtid + ,exchange_base_url + ,metadata + ,amount + ,credit_payto + ,initiation_time + ,status + ,status_msg + FROM transfer_operations + JOIN initiated_outgoing_transactions USING (initiated_outgoing_transaction_id) + WHERE initiated_outgoing_transaction_id=$1 + ", + ) + .bind(id as i64) + .try_map(|r: PgRow| { + Ok(TransferStatus { + status: r + .try_get::<SubmissionState, _>("status")? + .to_transfer_status(), + status_msg: r.try_get("status_msg")?, + amount: r.try_get_amount("amount", currency)?, + origin_exchange_url: r.try_get("exchange_base_url")?, + metadata: r.try_get("metadata")?, + wtid: r.try_get("wtid")?, + credit_account: r.try_get_payto("credit_payto")?, + timestamp: r.try_get_timestamp("initiation_time")?.into(), + }) + }) + .fetch_optional(db) + ) +} + +pub async fn transfer_page( + db: &PgPool, + currency: &Currency, + params: &Page, + status: &Option<TransferState>, +) -> sqlx::Result<Vec<TransferListStatus>> { + page( + db, + "initiated_outgoing_transaction_id", + params, + || { + let mut builder = QueryBuilder::new( + " + SELECT + initiated_outgoing_transaction_id + ,amount + ,status + ,credit_payto + ,initiation_time + FROM transfer_operations + JOIN initiated_outgoing_transactions USING (initiated_outgoing_transaction_id) + WHERE + ", + ); + if let Some(status) = status { + match status { + TransferState::pending => { + builder.push("( status = ").push_bind(SubmissionState::pending).push(" OR ").push(" status = ").push_bind(SubmissionState::unsubmitted).push(") AND "); + } + status => { + builder.push(" status = ").push_bind(SubmissionState::from(*status)).push(" AND ");} + } + } + builder + }, + |r: PgRow| { + Ok(TransferListStatus { + row_id: r.try_get_safeu64("initiated_outgoing_transaction_id")?, + status: r.try_get::<SubmissionState, _>("status")?.to_transfer_status(), + amount: r.try_get_amount("amount", currency)?, + credit_account: r.try_get_payto("credit_payto")?, + timestamp: r.try_get_timestamp("initiation_time")?.into(), + }) + }, + ) + .await +} diff --git a/src/db/payment.rs b/src/db/payment.rs @@ -335,30 +335,28 @@ mod test { use jiff::Timestamp; use sqlx::{PgPool, postgres::PgRow}; - use taler_api::{ - db::{IncomingType, TypeHelper as _}, - subject::subject_fmt_qr_bill, - }; + use taler_api::{db::TypeHelper as _, subject::subject_fmt_qr_bill}; use taler_common::{ api_common::{EddsaPublicKey, EddsaSignature, ShortHashCode}, + db::IncomingType, types::amount::amount, }; + use taler_test_utils::routine::Status::*; use uuid::Uuid; use crate::{ config::{AccountType, NexusIngestConfig}, db::{ - RegistrationResult, initiated::{PaymentInitiationResult, batch_initiated, initiate, initiated_ack}, payment::{ InResult, IncomingBounceRegistrationResult, OutgoingRegistrationResult, register_in_malformed, }, test::{ - CURRENCY, Status::*, check_in, check_in_count, check_out_count, gen_in_pay, + CURRENCY, check_in_count, check_in_state, check_out_count, gen_in_pay, gen_init_pay, gen_out_pay, setup, }, - transfer_register, + transfer::{RegistrationResult, transfer_register}, }, model::{IncomingId, IncomingPayment, OutgoingBatch, OutgoingId, OutgoingPayment}, rand_ebics_id, @@ -648,24 +646,24 @@ mod test { // Register let incoming = gen_in_pay("test".to_owned()); register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Bounced]).await; + check_in_state(&db, &[Bounced]).await; // Idempotent register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Bounced]).await; + check_in_state(&db, &[Bounced]).await; // Many register_incoming(&db, &cfg, &gen_in_pay("another subject".to_owned())) .await .unwrap(); - check_in(&db, &[Bounced, Bounced]).await; + check_in_state(&db, &[Bounced, Bounced]).await; // Admin balance adjust is ignored register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned())) .await .unwrap(); - check_in(&db, &[Bounced, Bounced, Simple]).await; + check_in_state(&db, &[Bounced, Bounced, Simple]).await; let original = gen_in_pay("test 2".to_owned()); let incomplete = IncomingPayment { @@ -676,13 +674,13 @@ mod test { // Register incomplete transaction register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await; + check_in_state(&db, &[Bounced, Bounced, Simple, Incomplete]).await; // Idempotent register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await; + check_in_state(&db, &[Bounced, Bounced, Simple, Incomplete]).await; // Recover info when completed register_incoming(&db, &cfg, &original).await.unwrap(); - check_in(&db, &[Bounced, Bounced, Simple, Bounced]).await; + check_in_state(&db, &[Bounced, Bounced, Simple, Bounced]).await; } #[tokio::test] @@ -696,11 +694,11 @@ mod test { // Register let incoming = gen_in_pay(subject.clone()); register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Reserve(key.clone())]).await; + check_in_state(&db, &[Reserve(key.clone())]).await; // Idempotent register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Reserve(key.clone())]).await; + check_in_state(&db, &[Reserve(key.clone())]).await; // Key reuse is bounced register_incoming(&db, &cfg, &gen_in_pay(subject.clone())) @@ -709,13 +707,13 @@ mod test { register_incoming(&db, &cfg, &gen_in_pay(format!("another {subject}"))) .await .unwrap(); - check_in(&db, &[Reserve(key.clone()), Bounced, Bounced]).await; + check_in_state(&db, &[Reserve(key.clone()), Bounced, Bounced]).await; // Admin balance adjust is ignored register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned())) .await .unwrap(); - check_in(&db, &[Reserve(key.clone()), Bounced, Bounced, Simple]).await; + check_in_state(&db, &[Reserve(key.clone()), Bounced, Bounced, Simple]).await; let new = EddsaPublicKey::rand(); let original = gen_in_pay(format!("test 2 with {new} reserve pub")); @@ -727,21 +725,21 @@ mod test { // Register incomplete transaction register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in( + check_in_state( &db, &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete], ) .await; // Idempotent register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in( + check_in_state( &db, &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete], ) .await; // Recover info when completed register_incoming(&db, &cfg, &original).await.unwrap(); - check_in( + check_in_state( &db, &[Reserve(key.clone()), Bounced, Bounced, Simple, Reserve(new)], ) @@ -777,17 +775,17 @@ mod test { // Register let incoming = gen_in_pay(subject.clone()); register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Reserve(first.clone())]).await; + check_in_state(&db, &[Reserve(first.clone())]).await; // Idempotent register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Reserve(first.clone())]).await; + check_in_state(&db, &[Reserve(first.clone())]).await; // Admin balance adjust is ignored register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned())) .await .unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple]).await; let original = gen_in_pay(format!("test 2 for {subject}")); let incomplete = IncomingPayment { @@ -797,13 +795,13 @@ mod test { }; // Register incomplete transaction register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; // Idempotent register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; // Recover info when completed register_incoming(&db, &cfg, &original).await.unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Bounced]).await; let second = EddsaPublicKey::rand(); assert_eq!( @@ -821,7 +819,7 @@ mod test { .unwrap(), RegistrationResult::Success ); - check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Bounced]).await; // Key reuse is pending for _ in 0..3 { @@ -829,7 +827,7 @@ mod test { .await .unwrap(); } - check_in( + check_in_state( &db, &[ Reserve(first.clone()), @@ -859,7 +857,7 @@ mod test { .unwrap(), RegistrationResult::Success ); - check_in( + check_in_state( &db, &[ Reserve(first.clone()), @@ -901,17 +899,17 @@ mod test { // Register let incoming = gen_in_pay(reference_number.clone()); register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Reserve(first.clone())]).await; + check_in_state(&db, &[Reserve(first.clone())]).await; // Idempotent register_incoming(&db, &cfg, &incoming).await.unwrap(); - check_in(&db, &[Reserve(first.clone())]).await; + check_in_state(&db, &[Reserve(first.clone())]).await; // Admin balance adjust is ignored register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned())) .await .unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple]).await; let original = gen_in_pay(reference_number.clone()); let incomplete = IncomingPayment { @@ -921,13 +919,13 @@ mod test { }; // Register incomplete transaction register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; // Idempotent register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Incomplete]).await; // Recover info when completed register_incoming(&db, &cfg, &original).await.unwrap(); - check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Bounced]).await; let second = EddsaPublicKey::rand(); assert_eq!( @@ -945,7 +943,7 @@ mod test { .unwrap(), RegistrationResult::Success ); - check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await; + check_in_state(&db, &[Reserve(first.clone()), Simple, Bounced]).await; // Key reuse is pending for _ in 0..3 { @@ -953,7 +951,7 @@ mod test { .await .unwrap(); } - check_in( + check_in_state( &db, &[ Reserve(first.clone()), @@ -983,7 +981,7 @@ mod test { .unwrap(), RegistrationResult::Success ); - check_in( + check_in_state( &db, &[ Reserve(first.clone()), @@ -1151,7 +1149,7 @@ mod test { register_incoming(&db, &cfg, &incomplete).await.unwrap(); register_incoming(&db, &cfg, &payment).await.unwrap(); register_incoming(&db, &cfg, &incomplete).await.unwrap(); - check_in(&db, &[Reserve(key.clone())]).await; + check_in_state(&db, &[Reserve(key.clone())]).await; // Check we do not register as talerable bounced transaction let new_key = EddsaPublicKey::rand(); @@ -1164,6 +1162,6 @@ mod test { register_incoming(&db, &cfg, &payment).await.unwrap(); register_incoming(&db, &cfg, &incomplete).await.unwrap(); register_incoming(&db, &cfg, &payment).await.unwrap(); - check_in(&db, &[Reserve(key.clone()), Bounced]).await; + check_in_state(&db, &[Reserve(key.clone()), Bounced]).await; } } diff --git a/src/db/transfer.rs b/src/db/transfer.rs @@ -0,0 +1,83 @@ +/* + This file is part of TALER + Copyright (C) 2026 Taler Systems SA + + TALER is free software; you can redistribute it and/or modify it under the + terms of the GNU Affero General Public License as published by the Free Software + Foundation; either version 3, or (at your option) any later version. + + TALER is distributed in the hope that it will be useful, but WITHOUT ANY + WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR + A PARTICULAR PURPOSE. See the GNU Affero General Public License for more details. + + You should have received a copy of the GNU Affero General Public License along with + TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> +*/ + +use jiff::Timestamp; +use sqlx::{PgPool, Row as _, postgres::PgRow}; +use taler_api::db::{BindHelper as _, TypeHelper as _}; +use taler_common::{ + api_common::{EddsaPublicKey, EddsaSignature}, + db::IncomingType, +}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RegistrationResult { + Success, + ReservePubReuse, + SubjectReuse, +} + +pub async fn transfer_register( + db: &PgPool, + ty: IncomingType, + account_pub: &EddsaPublicKey, + auth_pub: &EddsaPublicKey, + auth_sig: &EddsaSignature, + recurrent: bool, + reference_number: &str, + timestamp: &Timestamp, +) -> sqlx::Result<RegistrationResult> { + sqlx::query( + " + SELECT + out_subject_reuse, + out_reserve_pub_reuse + FROM register_prepared_transfers ( + $1::taler_incoming_type,$2,$3,$4,$5,$6,$7 + ) + ", + ) + .bind(ty) + .bind(account_pub) + .bind(auth_pub) + .bind(auth_sig) + .bind(recurrent) + .bind(reference_number) + .bind_timestamp(timestamp) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_subject_reuse")? { + RegistrationResult::SubjectReuse + } else if r.try_get_flag("out_reserve_pub_reuse")? { + RegistrationResult::ReservePubReuse + } else { + RegistrationResult::Success + }) + }) + .fetch_one(db) + .await +} + +pub async fn transfer_unregister( + db: &PgPool, + auth_pub: &EddsaPublicKey, + timestamp: &Timestamp, +) -> sqlx::Result<bool> { + sqlx::query("SELECT out_found FROM delete_prepared_transfers($1,$2)") + .bind(auth_pub) + .bind_timestamp(timestamp) + .try_map(|r: PgRow| r.try_get(0)) + .fetch_one(db) + .await +} diff --git a/src/lib.rs b/src/lib.rs @@ -39,6 +39,7 @@ use crate::{ }; use crate::{keys::ClientPriKeysFile, xml::XmlReader}; +pub mod api; pub mod common; pub mod config; pub mod crypto; diff --git a/src/model.rs b/src/model.rs @@ -47,7 +47,7 @@ pub enum SubmissionState { } impl SubmissionState { - fn to_transfer_status(self) -> TransferState { + pub fn to_transfer_status(self) -> TransferState { match self { SubmissionState::unsubmitted | SubmissionState::pending => TransferState::pending, SubmissionState::transient_failure => TransferState::transient_failure, @@ -57,6 +57,18 @@ impl SubmissionState { } } +impl From<TransferState> for SubmissionState { + fn from(value: TransferState) -> Self { + match value { + TransferState::pending => SubmissionState::pending, + TransferState::transient_failure => SubmissionState::transient_failure, + TransferState::permanent_failure => SubmissionState::permanent_failure, + TransferState::late_failure => SubmissionState::late_failure, + TransferState::success => SubmissionState::success, + } + } +} + /// ID for incoming transactions #[derive(Debug, Clone)] pub struct IncomingId {