libeufin

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

commit 895ca7f0a1d791292910db8f593b6dd2a5938229
parent efe0ff0ffba560e68bb4ef9e87119076e79d9f5b
Author: Antoine A <>
Date:   Fri, 22 May 2026 10:50:01 +0200

common: handle serialization in all database query and clean code

Diffstat:
MCargo.toml | 7+++++++
MMakefile | 12++++++++++++
Mcrates/libeufin-bank/src/api.rs | 28++++++++++++++--------------
Mcrates/libeufin-bank/src/api/account.rs | 13+++++++------
Mcrates/libeufin-bank/src/api/prepared.rs | 2+-
Mcrates/libeufin-bank/src/api/tan.rs | 34+++++++++++++++++-----------------
Mcrates/libeufin-bank/src/api/token.rs | 16++++++++--------
Mcrates/libeufin-bank/src/api/tx.rs | 2+-
Mcrates/libeufin-bank/src/db/account.rs | 717++++++++++++++++++++++++++++++++++++++++---------------------------------------
Mcrates/libeufin-bank/src/db/exchange.rs | 4++--
Mcrates/libeufin-bank/src/db/prepared.rs | 2+-
Mcrates/libeufin-bank/src/db/tan.rs | 141++++++++++++++++++++++++++++++++++++++++++-------------------------------------
Mcrates/libeufin-bank/src/db/tx.rs | 83++++++++++++++++++++++++++++++++++++++++---------------------------------------
Mcrates/libeufin-ebics/src/db.rs | 8++++----
Mcrates/libeufin-ebics/src/iso20022/pain001.rs | 6+++---
Mcrates/libeufin-ebics/src/ws.rs | 4++--
Mcrates/libeufin-nexus/src/bench.rs | 12++++++------
Mcrates/libeufin-nexus/src/db.rs | 2+-
Mcrates/libeufin-nexus/src/db/exchange.rs | 2+-
Mcrates/libeufin-nexus/src/db/initiated.rs | 463+++++++++++++++++++++++++++++++++++++++++--------------------------------------
Mcrates/libeufin-nexus/src/db/list.rs | 172++++++++++++++++++++++++++++++++++++++++---------------------------------------
Mcrates/libeufin-nexus/src/db/payment.rs | 311++++++++++++++++++++++++++++++++++++++++---------------------------------------
Mcrates/libeufin-nexus/src/db/transfer.rs | 2+-
Mcrates/libeufin-nexus/src/test.rs | 7+++----
24 files changed, 1064 insertions(+), 986 deletions(-)

diff --git a/Cargo.toml b/Cargo.toml @@ -10,6 +10,13 @@ homepage = "https://taler.net/" repository = "https://git.taler.net/libeufin.git" license-file = "COPYING" +[profile.dev] +debug = true + +[profile.release] +lto = "fat" +codegen-units = 1 + [workspace.dependencies] axum = { version = "0.8", features = ["ws", "macros"] } tracing = "0.1" diff --git a/Makefile b/Makefile @@ -130,3 +130,15 @@ bank-bench-db: install-nobuild-files .PHONY: nexus-bench-db nexus-bench-db: install-nobuild-files ./gradlew cleanTest :libeufin-nexus:test --tests Bench.benchDb -i --no-build-cache + +.PHONY: rust-install +rust-install: install-nobuild-files + cargo build --release --bin libeufin-bank --bin libeufin-rust + install -D -t $(bin_dir) target/release/libeufin-bank + install -D -t $(bin_dir) target/release/libeufin-rust + +.PHONY: rust-check +rust-check: install-nobuild-files + cargo clippy --all-targets + cargo test + diff --git a/crates/libeufin-bank/src/api.rs b/crates/libeufin-bank/src/api.rs @@ -263,7 +263,7 @@ pub mod test { None, None, None, - rand_iban_payto().into_inner().into(), + &rand_iban_payto().into_inner().into(), false, false, Decimal::new(10, 0).to_amount(&state.cfg.regional_currency), @@ -285,7 +285,7 @@ pub mod test { None, None, None, - rand_iban_payto().into_inner().into(), + &rand_iban_payto().into_inner().into(), false, true, Decimal::new(10, 0).to_amount(&state.cfg.regional_currency), @@ -307,7 +307,7 @@ pub mod test { None, None, None, - rand_iban_payto().into_inner().into(), + &rand_iban_payto().into_inner().into(), false, false, Decimal::new(10, 0).to_amount(&state.cfg.regional_currency), @@ -390,7 +390,7 @@ pub mod test { fn extract_username(path: &str) -> &str { if path.contains("admin") { - return "admin"; + "admin" } else { path.split('/').nth(2).unwrap() } @@ -398,12 +398,12 @@ pub mod test { pub async fn cache_tokens(&mut self, usernames: &[&'static str]) { let tasks = usernames - .into_iter() + .iter() .map(|username| { let username = CompactString::from(*username); async { let res = Self::pw_auth( - self.server.post(&format!("/accounts/{username}/token")), + self.server.post(format!("/accounts/{username}/token")), Some("admin"), ) .json(json!({ @@ -537,7 +537,7 @@ pub mod test { } pub async fn fill_tan_info(&self, username: &str) { - self.patch_admin(&format!("/accounts/{username}")) + self.patch_admin(format!("/accounts/{username}")) .json(json!({ "contact_data": { "phone": format!("+{}", random_range(0..10000)) @@ -549,7 +549,7 @@ pub mod test { } pub async fn fill_cashout_info(&self, username: &str) { - self.patch_admin(&format!("/accounts/{username}")) + self.patch_admin(format!("/accounts/{username}")) .json(json!({ "cashout_payto_uri": self.unknown_payto, })) @@ -576,7 +576,7 @@ pub mod test { "amount": format!("{}:{amount}", self.state.cfg.regional_currency), })) .await - .maybe_challenge(&self) + .maybe_challenge(self) .await .assert_ok(); } @@ -670,7 +670,7 @@ pub mod test { /** Set [account] debit threshold to [maxDebt] amount */ pub async fn set_max_debt(&self, username: &str, amount: &str) { - self.patch_admin(&format!("/accounts/{username}")) + self.patch_admin(format!("/accounts/{username}")) .json(json!({ "debit_threshold": format!("{}:{amount}", self.state.cfg.regional_currency) })) @@ -681,7 +681,7 @@ pub mod test { /** Check [account] balance is [amount], [amount] is prefixed with + for credit and - for debit */ pub async fn assert_balance(&self, username: &str, expected: &str) { let res: AccountData = self - .get_admin(&format!("/accounts/{username}")) + .get_admin(format!("/accounts/{username}")) .await .assert_ok_json(); let Balance { @@ -722,7 +722,7 @@ pub mod test { let code = match std::fs::read_to_string(&path) { Ok(f) => f, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return None, - Err(e) => Err(e).unwrap(), + Err(e) => panic!("{:?}", e), }; std::fs::remove_file(path).unwrap(); Some(code.split(' ').next().unwrap().into()) @@ -758,7 +758,7 @@ pub mod test { }; for challenge in challenges { - ctx.posta(&format!( + ctx.posta(format!( "/accounts/{username}/challenge/{}", challenge.challenge_id )) @@ -769,7 +769,7 @@ pub mod test { for challenge in challenges { let code = tan_code(&challenge.tan_info).unwrap(); - ctx.posta(&format!( + ctx.posta(format!( "/accounts/{username}/challenge/{}/confirm", challenge.challenge_id )) diff --git a/crates/libeufin-bank/src/api/account.rs b/crates/libeufin-bank/src/api/account.rs @@ -528,6 +528,7 @@ pub fn account_api() -> Router<Arc<BankState>> { .await? { PatchResult::Success => Ok(NoContent.into_response()), + PatchResult::MissingTanInfo(e) => Err(e), PatchResult::Challenges(tans) => { if tans.is_empty() { mfa.response_mfa(&mut auth, &state.db, &state.cfg.ctx).await @@ -727,7 +728,7 @@ pub async fn create_account( .and_then(|it| it.phone.opt()) .map(|it| it.as_str()), req.cashout_payto_uri.as_ref(), - payto, + &payto, req.is_public, req.is_taler_exchange, req.debit_threshold.unwrap_or(cfg.default_debt_limit), @@ -812,7 +813,7 @@ pub async fn create_admin_account( None, None, None, - payto, + &payto, false, false, cfg.default_debt_limit, @@ -863,7 +864,7 @@ pub async fn patch_account( } } - reconfig( + Ok(reconfig( db, &cfg.regional_currency, username, @@ -873,7 +874,7 @@ pub async fn patch_account( cfg.allow_edit_name, cfg.allow_edit_cashout, ) - .await + .await?) } #[cfg(test)] @@ -1639,13 +1640,13 @@ pub mod test { .await .assert_challenge_check(&ctx, async |_| { let acc: AccountData = ctx.geta("/accounts/customer").await.assert_ok_json(); - assert_eq!(acc.is_public, false); + assert!(!acc.is_public); }) .await .assert_no_content(); let acc: AccountData = ctx.geta("/accounts/customer").await.assert_ok_json(); - assert_eq!(acc.is_public, true); + assert!(acc.is_public); // Restriction let ctx = ctx.swap_cfg("test_restrict.conf").await; diff --git a/crates/libeufin-bank/src/api/prepared.rs b/crates/libeufin-bank/src/api/prepared.rs @@ -189,7 +189,7 @@ mod test { } else if r.try_get_flag(1)? { Status::Bounced } else { - match r.try_get(2)? { + match r.try_get_opt_parse(2)? { None => Status::Simple, Some(IncomingType::reserve) => Status::Reserve(r.try_get(3)?), Some(IncomingType::kyc) => Status::Kyc(r.try_get(3)?), diff --git a/crates/libeufin-bank/src/api/tan.rs b/crates/libeufin-bank/src/api/tan.rs @@ -200,7 +200,7 @@ pub mod test { }; let send = async |c: &Challenge| { - ctx.posta(&format!("/accounts/merchant/challenge/{}", c.challenge_id)) + ctx.posta(format!("/accounts/merchant/challenge/{}", c.challenge_id)) .await .assert_ok(); }; @@ -384,7 +384,7 @@ pub mod test { .assert_no_content(); // Check invalidated - ctx.posta(&format!( + ctx.posta(format!( "/accounts/merchant/challenge/{}/confirm", challenge.challenge_id, )) @@ -401,7 +401,7 @@ pub mod test { } // Unknown challenge - ctx.posta(&format!("/accounts/merchant/challenge/{}", Uuid::new_v4())) + ctx.posta(format!("/accounts/merchant/challenge/{}", Uuid::new_v4())) .await .assert_error(ErrorCode::BANK_TRANSACTION_NOT_FOUND); @@ -430,7 +430,7 @@ pub mod test { }; let submit = async |c: &Challenge| { - ctx.posta(&format!("/accounts/merchant/challenge/{}", c.challenge_id)) + ctx.posta(format!("/accounts/merchant/challenge/{}", c.challenge_id)) .await .assert_ok_json::<ChallengeRequestResponse>() }; @@ -450,13 +450,13 @@ pub mod test { submit(&tx_challenge().await).await; } let c = tx_challenge().await; - ctx.posta(&format!("/accounts/merchant/challenge/{}", c.challenge_id)) + ctx.posta(format!("/accounts/merchant/challenge/{}", c.challenge_id)) .await .assert_error(ErrorCode::BANK_TAN_RATE_LIMITED); // Old already submitted challenge still works submit(&old).await; - ctx.posta(&format!( + ctx.posta(format!( "/accounts/merchant/challenge/{}/confirm", old.challenge_id )) @@ -469,7 +469,7 @@ pub mod test { // We are still rate limited let new = tx_challenge().await; - ctx.posta(&format!( + ctx.posta(format!( "/accounts/merchant/challenge/{}", new.challenge_id )) @@ -489,7 +489,7 @@ pub mod test { .await .assert_accepted_json(); let challenge = &res.challenges[0]; - ctx.posta(&format!( + ctx.posta(format!( "/accounts/merchant/challenge/{}", challenge.challenge_id )) @@ -511,36 +511,36 @@ pub mod test { .assert_accepted_json(); let challenge = &res.challenges[0]; let id = &challenge.challenge_id; - ctx.posta(&format!("/accounts/merchant/challenge/{id}")) + ctx.posta(format!("/accounts/merchant/challenge/{id}")) .await .assert_ok_json::<ChallengeRequestResponse>(); let code = tan_code(&challenge.tan_info); // Check bad TAN code - ctx.posta(&format!("/accounts/merchant/challenge/{id}/confirm")) + ctx.posta(format!("/accounts/merchant/challenge/{id}/confirm")) .json(json!({ "tan": "nice-try" })) .await .assert_error(ErrorCode::BANK_TAN_CHALLENGE_FAILED); // Check wrong account - ctx.posta(&format!("/accounts/customer/challenge/{id}/confirm")) + ctx.posta(format!("/accounts/customer/challenge/{id}/confirm")) .json(json!({ "tan": "nice-try" })) .await .assert_error(ErrorCode::BANK_TAN_CHALLENGE_FAILED); // Check OK - ctx.posta(&format!("/accounts/customer/challenge/{id}/confirm")) + ctx.posta(format!("/accounts/customer/challenge/{id}/confirm")) .json(json!({ "tan": code })) .await .assert_no_content(); // Check idempotence - ctx.posta(&format!("/accounts/customer/challenge/{id}/confirm")) + ctx.posta(format!("/accounts/customer/challenge/{id}/confirm")) .json(json!({ "tan": code })) .await .assert_no_content(); // Unknown challenge - ctx.posta(&format!( + ctx.posta(format!( "/accounts/customer/challenge/{}/confirm", Uuid::new_v4() )) @@ -558,18 +558,18 @@ pub mod test { .assert_accepted_json(); let challenge = &res.challenges[0]; let id = &challenge.challenge_id; - ctx.posta(&format!("/accounts/merchant/challenge/{id}")) + ctx.posta(format!("/accounts/merchant/challenge/{id}")) .await .assert_ok_json::<ChallengeRequestResponse>(); // Check invalidated ctx.fill_tan_info("merchant").await; - ctx.posta(&format!("/accounts/customer/challenge/{id}/confirm")) + ctx.posta(format!("/accounts/customer/challenge/{id}/confirm")) .json(json!({ "tan": tan_code(&challenge.tan_info) })) .await .assert_error(ErrorCode::BANK_TAN_CHALLENGE_EXPIRED); - ctx.posta(&format!("/accounts/customer/challenge/{id}")) + ctx.posta(format!("/accounts/customer/challenge/{id}")) .await .assert_error(ErrorCode::BANK_TAN_CHALLENGE_EXPIRED); } diff --git a/crates/libeufin-bank/src/api/token.rs b/crates/libeufin-bank/src/api/token.rs @@ -391,14 +391,14 @@ pub mod test { .await .assert_ok_json(); // Wrong account - ctx.deletea(&format!("/accounts/customer/tokens/{}", res.token_id)) + ctx.deletea(format!("/accounts/customer/tokens/{}", res.token_id)) .await .assert_error(ErrorCode::BANK_TRANSACTION_NOT_FOUND); // Check OK - ctx.deletea(&format!("/accounts/merchant/tokens/{}", res.token_id)) + ctx.deletea(format!("/accounts/merchant/tokens/{}", res.token_id)) .await .assert_no_content(); - ctx.deletea(&format!("/accounts/merchant/tokens/{}", res.token_id)) + ctx.deletea(format!("/accounts/merchant/tokens/{}", res.token_id)) .await .assert_error(ErrorCode::BANK_TRANSACTION_NOT_FOUND); // Check token no longer work @@ -429,7 +429,7 @@ pub mod test { .challenges .pop() .unwrap(); - ctx.post(&format!( + ctx.post(format!( "/accounts/merchant/challenge/{}", challenge.challenge_id )) @@ -437,7 +437,7 @@ pub mod test { .assert_ok(); assert_eq!("REDACTED", challenge.tan_info); // Check phone number is hidden let code = tan_code("+12345"); - ctx.post(&format!( + ctx.post(format!( "/accounts/merchant/challenge/{}/confirm", challenge.challenge_id )) @@ -464,7 +464,7 @@ pub mod test { .challenges .pop() .unwrap(); - ctx.post(&format!( + ctx.post(format!( "/accounts/merchant/challenge/{}", challenge.challenge_id )) @@ -472,7 +472,7 @@ pub mod test { .assert_ok(); while counter > 0 { let error: ErrorDetail = ctx - .post(&format!( + .post(format!( "/accounts/merchant/challenge/{}/confirm", challenge.challenge_id )) @@ -512,7 +512,7 @@ pub mod test { // Check OK for account in ["merchant", "customer"] { - ctx.geta(&format!("/accounts/{account}/tokens")) + ctx.geta(format!("/accounts/{account}/tokens")) .await .assert_no_content(); } diff --git a/crates/libeufin-bank/src/api/tx.rs b/crates/libeufin-bank/src/api/tx.rs @@ -111,7 +111,7 @@ pub fn tx_api() -> Router<Arc<BankState>> { amount, &Timestamp::now(), mfa.is_2fa(), - req.request_uid, + &req.request_uid, state.cfg.wire_transfer_fees, state.cfg.min_amount, state.cfg.max_amount, diff --git a/crates/libeufin-bank/src/db/account.rs b/crates/libeufin-bank/src/db/account.rs @@ -21,8 +21,8 @@ use sqlx::{ postgres::{PgArguments, PgRow}, }; use taler_api::{ - db::{BindHelper as _, PgError, TypeHelper as _, page}, - error::ApiResult, + db::{BindHelper as _, TypeHelper as _, page}, + error::ApiError, serialized, }; use taler_common::{ @@ -79,7 +79,7 @@ pub async fn create( email: Option<&str>, phone: Option<&str>, cashout: Option<&IbanPayto>, - internal: LibeufinId, + internal: &LibeufinId, is_public: bool, is_exchange: bool, max_debt: Amount, @@ -88,13 +88,13 @@ pub async fn create( check_payto_idempotent: bool, conversion_rate_class_id: Option<u64>, ) -> sqlx::Result<CreationResult> { - // TODO serialized - let mut tx = db.begin().await?; - let now = Timestamp::now(); - let cashout = cashout.map(|it| it.to_string()); - let canonical = internal.canonical(); - let idempotent = sqlx::query( - " + let canonical = &internal.canonical(); + serialized!(async { + let mut tx = db.begin().await?; + let now = Timestamp::now(); + let cashout = cashout.map(|it| it.to_string()); + let idempotent = sqlx::query( + " SELECT password_hash, name=$1 AND email IS NOT DISTINCT FROM $2 AND phone IS NOT DISTINCT FROM $3 @@ -111,51 +111,52 @@ pub async fn create( ON customer_id=owning_customer_id WHERE username=$12 ", - ) - .bind(name) - .bind(email) - .bind(phone) - .bind(&cashout) - .bind(tan_channels) - .bind(check_payto_idempotent) - .bind(&canonical) - .bind(is_public) - .bind(is_exchange) - .bind(max_debt) - .bind(conversion_rate_class_id.map(|it| it as i64)) - .bind(username) - .try_map(|r: PgRow| { - Ok(( - r.try_get(1)? && pw_crypto.checkpw(password, r.try_get(0)?).unwrap().matches, - sql_bank_payto(&r, ctx, "internal_payto", "name")?, - )) - }) - .fetch_optional(&mut *tx) - .await?; - let res = if let Some((matches, payto)) = idempotent { - if matches { - CreationResult::Success(payto) + ) + .bind(name) + .bind(email) + .bind(phone) + .bind(&cashout) + .bind(tan_channels) + .bind(check_payto_idempotent) + .bind(canonical) + .bind(is_public) + .bind(is_exchange) + .bind(max_debt) + .bind(conversion_rate_class_id.map(|it| it as i64)) + .bind(username) + .try_map(|r: PgRow| { + Ok(( + r.try_get(1)? && pw_crypto.checkpw(password, r.try_get(0)?).unwrap().matches, + sql_bank_payto(&r, ctx, "internal_payto", "name")?, + )) + }) + .fetch_optional(&mut *tx) + .await?; + let res = if let Some((matches, payto)) = idempotent { + if matches { + CreationResult::Success(payto) + } else { + CreationResult::UsernameReuse + } } else { - CreationResult::UsernameReuse - } - } else { - if let LibeufinId::IBAN(BankID { iban, .. }) = &internal { - let res = sqlx::query("INSERT INTO iban_history(iban,creation_time) VALUES ($1, $2)") - .bind(iban.as_ref()) - .bind_timestamp(&now) - .execute(&mut *tx) - .await; - if let Err(e) = &res - && e.is_unique_err() - { - tx.rollback().await?; - return sqlx::Result::Ok(CreationResult::PayToReuse); + if let LibeufinId::IBAN(BankID { iban, .. }) = &internal { + let res = + sqlx::query("INSERT INTO iban_history(iban,creation_time) VALUES ($1, $2)") + .bind(iban.as_ref()) + .bind_timestamp(&now) + .execute(&mut *tx) + .await; + if let Err(e) = &res + && e.is_unique_err() + { + tx.rollback().await?; + return Ok(CreationResult::PayToReuse); + } + res?; } - res?; - } - let customer_id: i64 = sqlx::query_scalar( - " + let customer_id: i64 = sqlx::query_scalar( + " INSERT INTO customers ( username ,password_hash @@ -167,18 +168,18 @@ pub async fn create( ) VALUES ($1, $2, $3, $4, $5, $6, sort_uniq($7)) RETURNING customer_id ", - ) - .bind(username) - .bind(pw_crypto.hashpw(password)) - .bind(name) - .bind(email) - .bind(phone) - .bind(&cashout) - .bind(tan_channels) - .fetch_one(&mut *tx) - .await?; - let res = sqlx::query( - " + ) + .bind(username) + .bind(pw_crypto.hashpw(password)) + .bind(name) + .bind(email) + .bind(phone) + .bind(&cashout) + .bind(tan_channels) + .fetch_one(&mut *tx) + .await?; + let res = sqlx::query( + " INSERT INTO bank_accounts( internal_payto ,owning_customer_id @@ -188,44 +189,45 @@ pub async fn create( ,conversion_rate_class_id ) VALUES ($1, $2, $3, $4, $5, $6) ", - ) - .bind(&canonical) - .bind(customer_id) - .bind(is_public) - .bind(is_exchange) - .bind(max_debt) - .bind(conversion_rate_class_id.map(|it| it as i64)) - .execute(&mut *tx) - .await; + ) + .bind(canonical) + .bind(customer_id) + .bind(is_public) + .bind(is_exchange) + .bind(max_debt) + .bind(conversion_rate_class_id.map(|it| it as i64)) + .execute(&mut *tx) + .await; - if let Err(e) = &res - && e.is_unique_err() - { - tx.rollback().await?; - return sqlx::Result::Ok(CreationResult::PayToReuse); - } else if let Err(e) = &res - && e.is_fk_err() - { - tx.rollback().await?; - return sqlx::Result::Ok(CreationResult::UnknownConversionClass); - } - res?; + if let Err(e) = &res + && e.is_unique_err() + { + tx.rollback().await?; + return Ok(CreationResult::PayToReuse); + } else if let Err(e) = &res + && e.is_fk_err() + { + tx.rollback().await?; + return Ok(CreationResult::UnknownConversionClass); + } + res?; - if !bonus.is_zero() { - let insufisient = sqlx::query_scalar(" + if !bonus.is_zero() { + let insufisient = sqlx::query_scalar(" SELECT out_balance_insufficient FROM bank_transaction($1,'admin','bonus',$2,$3,true,NULL,NULL,NULL,NULL, NULL, NULL, NULL) - ").bind(&canonical).bind(bonus).bind(now.as_microsecond()).fetch_one(&mut *tx).await?; - if insufisient { - tx.rollback().await?; - return sqlx::Result::Ok(CreationResult::BonusBalanceInsufficient); + ").bind(canonical).bind(bonus).bind(now.as_microsecond()).fetch_one(&mut *tx).await?; + if insufisient { + tx.rollback().await?; + return Ok(CreationResult::BonusBalanceInsufficient); + } } - } - CreationResult::Success(internal.bank(ctx).into_inner().full(name)) - }; - tx.commit().await?; - sqlx::Result::Ok(res) + CreationResult::Success(internal.clone().bank(ctx).into_inner().full(name)) + }; + tx.commit().await?; + Ok(res) + }) } /** Result status of account deletion */ @@ -242,31 +244,32 @@ pub async fn delete( username: &str, is2fa: bool, ) -> sqlx::Result<DeletionResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_not_found, out_balance_not_zero, out_tan_required FROM account_delete($1,$2,$3) ", - ) - .bind(username) - .bind_timestamp(&Timestamp::now()) - .bind(is2fa) - .try_map(|r: PgRow| { - Ok(if r.try_get_flag("out_not_found")? { - DeletionResult::UnknownAccount - } else if r.try_get_flag("out_balance_not_zero")? { - DeletionResult::BalanceNotZero - } else if r.try_get_flag("out_tan_required")? { - DeletionResult::TanRequired - } else { - DeletionResult::Success + ) + .bind(username) + .bind_timestamp(&Timestamp::now()) + .bind(is2fa) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_not_found")? { + DeletionResult::UnknownAccount + } else if r.try_get_flag("out_balance_not_zero")? { + DeletionResult::BalanceNotZero + } else if r.try_get_flag("out_tan_required")? { + DeletionResult::TanRequired + } else { + DeletionResult::Success + }) }) - }) - .fetch_one(db) - .await + .fetch_one(db) + ) } /** Result status of customer account patch */ @@ -278,6 +281,7 @@ pub enum PatchResult { NonAdminConversionRateClass, UnknownConversionClass, Challenges(Tans), + MissingTanInfo(ApiError), Success, } @@ -291,7 +295,7 @@ pub async fn reconfig( is2fa: bool, allow_edit_name: bool, allow_edit_cashout: bool, -) -> ApiResult<PatchResult> { +) -> sqlx::Result<PatchResult> { let AccountReconfiguration { cashout_payto_uri, name, @@ -302,8 +306,6 @@ pub async fn reconfig( .. } = req; - let mut tx = db.begin().await?; - #[derive(Debug)] struct CurrentAccount { id: u64, @@ -329,10 +331,12 @@ pub async fn reconfig( &self.channels } } + serialized!(async { + let mut tx = db.begin().await?; - // Get user ID and current data - let curr = sqlx::query( - " + // Get user ID and current data + let curr = sqlx::query( + " SELECT customer_id, tan_channels, @@ -347,145 +351,150 @@ pub async fn reconfig( ON customer_id=owning_customer_id WHERE username=$1 AND deleted_at IS NULL ", - ) - .bind(username) - .try_map(|r: PgRow| { - Ok(CurrentAccount { - id: r.try_get_u64("customer_id")?, - channels: r.try_get("tan_channels")?, - email: r.try_get("email")?, - phone: r.try_get("phone")?, - name: r.try_get("name")?, - cashout_pay_to: r.try_get_opt_parse("cashout_payto")?, - debt_limit: r.try_get_amount("max_debt", currency)?, - conversion_rate_class_id: r.try_get_opt_u64("conversion_rate_class_id")?, + ) + .bind(username) + .try_map(|r: PgRow| { + Ok(CurrentAccount { + id: r.try_get_u64("customer_id")?, + channels: r.try_get("tan_channels")?, + email: r.try_get("email")?, + phone: r.try_get("phone")?, + name: r.try_get("name")?, + cashout_pay_to: r.try_get_opt_parse("cashout_payto")?, + debt_limit: r.try_get_amount("max_debt", currency)?, + conversion_rate_class_id: r.try_get_opt_u64("conversion_rate_class_id")?, + }) }) - }) - .fetch_optional(&mut *tx) - .await?; - let Some(curr) = curr else { - return Ok(PatchResult::UnknownAccount); - }; + .fetch_optional(&mut *tx) + .await?; + let Some(curr) = curr else { + return Ok(PatchResult::UnknownAccount); + }; - let validation = req.required_validation(&curr)?; - // Check performed 2fa check - if !is_admin && !is2fa { - // Check if mfa is required - if !curr.channels.is_empty() { - let mut tans = curr.mfa(); + let validation = match req.required_validation(&curr) { + Ok(v) => v, + Err(e) => return Ok(PatchResult::MissingTanInfo(e)) + }; - if tans.len() == 1 { - // Performs mfa and validation at the same time - tans.extend(validation); - return Ok(PatchResult::Challenges(tans)); - } else { - return Ok(PatchResult::Challenges(Vec::new())); + // Check performed 2fa check + if !is_admin && !is2fa { + // Check if mfa is required + if !curr.channels.is_empty() { + let mut tans = curr.mfa(); + + if tans.len() == 1 { + // Performs mfa and validation at the same time + tans.extend(validation); + return Ok(PatchResult::Challenges(tans)); + } else { + return Ok(PatchResult::Challenges(Vec::new())); + } } - } - // Check if validation is required - if !validation.is_empty() { - return Ok(PatchResult::Challenges(validation)); + // Check if validation is required + if !validation.is_empty() { + return Ok(PatchResult::Challenges(validation)); + } } - } - // Check reconfig rights - if !is_admin { - if !allow_edit_name - && let Some(name) = name - && name != curr.name - { - return Ok(PatchResult::NonAdminName); - } else if !allow_edit_cashout - && let Some(cashout) = cashout_payto_uri.inner() - && cashout != curr.cashout_pay_to.as_ref() - { - return Ok(PatchResult::NonAdminCashout); - } else if let Some(limit) = debit_threshold - && limit != &curr.debt_limit - { - return Ok(PatchResult::NonAdminDebtLimit); - } else if let Some(id) = conversion_rate_class_id.inner() - && id != curr.conversion_rate_class_id.as_ref() - { - return Ok(PatchResult::NonAdminConversionRateClass); + // Check reconfig rights + if !is_admin { + if !allow_edit_name + && let Some(name) = name + && name != curr.name + { + return Ok(PatchResult::NonAdminName); + } else if !allow_edit_cashout + && let Some(cashout) = cashout_payto_uri.inner() + && cashout != curr.cashout_pay_to.as_ref() + { + return Ok(PatchResult::NonAdminCashout); + } else if let Some(limit) = debit_threshold + && limit != &curr.debt_limit + { + return Ok(PatchResult::NonAdminDebtLimit); + } else if let Some(id) = conversion_rate_class_id.inner() + && id != curr.conversion_rate_class_id.as_ref() + { + return Ok(PatchResult::NonAdminConversionRateClass); + } } - } - // Update bank info - let mut sql = QueryBuilder::new("Update bank_accounts SET "); - let mut separated = sql.separated(','); - if let Some(v) = is_public { - separated.push("is_public=").push_bind_unseparated(v); - } - if let Some(v) = is_taler_exchange { - separated - .push("is_taler_exchange=") - .push_bind_unseparated(v); - } - if let Some(v) = debit_threshold { - separated.push("max_debt=").push_bind_unseparated(v); - } - if let Some(v) = conversion_rate_class_id.inner() { - separated - .push("conversion_rate_class_id=") - .push_bind_unseparated(v.map(|it| *it as i64)); - } - if !sql.sql().ends_with("SET ") { - sql.push(" WHERE owning_customer_id=") - .push_bind(curr.id as i64); - let res = sql.build().execute(&mut *tx).await; - if let Err(e) = &res - && e.is_fk_err() - { - tx.rollback().await?; - return Ok(PatchResult::UnknownConversionClass); + // Update bank info + let mut sql = QueryBuilder::new("Update bank_accounts SET "); + let mut separated = sql.separated(','); + if let Some(v) = is_public { + separated.push("is_public=").push_bind_unseparated(v); + } + if let Some(v) = is_taler_exchange { + separated + .push("is_taler_exchange=") + .push_bind_unseparated(v); + } + if let Some(v) = debit_threshold { + separated.push("max_debt=").push_bind_unseparated(v); + } + if let Some(v) = conversion_rate_class_id.inner() { + separated + .push("conversion_rate_class_id=") + .push_bind_unseparated(v.map(|it| *it as i64)); + } + if !sql.sql().ends_with("SET ") { + sql.push(" WHERE owning_customer_id=") + .push_bind(curr.id as i64); + let res = sql.build().execute(&mut *tx).await; + if let Err(e) = &res + && e.is_fk_err() + { + tx.rollback().await?; + return Ok(PatchResult::UnknownConversionClass); + } + res?; } - res?; - } - // Update customer info - let mut sql = QueryBuilder::new("UPDATE customers SET "); - let mut separated = sql.separated(','); - if let Some(v) = cashout_payto_uri.inner() { - separated - .push("cashout_payto=") - .push_bind_unseparated(v.map(|it| it.as_uri().to_string())); - } - if let Some(v) = req.contact_data.as_ref().and_then(|it| it.phone.inner()) { - separated.push("phone=").push_bind_unseparated(v); - } - if let Some(v) = req.contact_data.as_ref().and_then(|it| it.email.inner()) { - separated.push("email=").push_bind_unseparated(v); - } - if let Some(v) = req.channels() { - separated - .push("tan_channels=sort_uniq(") - .push_bind_unseparated(v) - .push_unseparated(')'); - } - if let Some(v) = &req.name { - separated.push("name=").push_bind_unseparated(v); - } - if !sql.sql().ends_with("SET ") { - sql.push(" WHERE customer_id=") - .push_bind(curr.id as i64) - .build() - .execute(&mut *tx) - .await?; - } + // Update customer info + let mut sql = QueryBuilder::new("UPDATE customers SET "); + let mut separated = sql.separated(','); + if let Some(v) = cashout_payto_uri.inner() { + separated + .push("cashout_payto=") + .push_bind_unseparated(v.map(|it| it.as_uri().to_string())); + } + if let Some(v) = req.contact_data.as_ref().and_then(|it| it.phone.inner()) { + separated.push("phone=").push_bind_unseparated(v); + } + if let Some(v) = req.contact_data.as_ref().and_then(|it| it.email.inner()) { + separated.push("email=").push_bind_unseparated(v); + } + if let Some(v) = req.channels() { + separated + .push("tan_channels=sort_uniq(") + .push_bind_unseparated(v) + .push_unseparated(')'); + } + if let Some(v) = &req.name { + separated.push("name=").push_bind_unseparated(v); + } + if !sql.sql().ends_with("SET ") { + sql.push(" WHERE customer_id=") + .push_bind(curr.id as i64) + .build() + .execute(&mut *tx) + .await?; + } - if !validation.is_empty() { - sqlx::query("UPDATE tan_challenges SET expiration_date=0 WHERE customer=$1") - .bind(curr.id as i64) - .execute(&mut *tx) - .await?; - } + if !validation.is_empty() { + sqlx::query("UPDATE tan_challenges SET expiration_date=0 WHERE customer=$1") + .bind(curr.id as i64) + .execute(&mut *tx) + .await?; + } - tx.commit().await?; + tx.commit().await?; - Ok(PatchResult::Success) + Ok(PatchResult::Success) + }) } /** Result status of customer account auth patch */ @@ -505,44 +514,46 @@ pub async fn reconfig_password( old_pw: Option<&str>, is2fa: bool, ) -> sqlx::Result<PatchAuthResult> { - // TODO use optimistic replace instead of transaction - let mut tx = db.begin().await?; + serialized!(async { + // TODO use optimistic replace instead of transaction + let mut tx = db.begin().await?; - let Some((customer_id, currenc_pwh, tan_required)): Option<(i64, String, bool)> = - sqlx::query_as( - " + let Some((customer_id, currenc_pwh, tan_required)): Option<(i64, String, bool)> = + sqlx::query_as( + " SELECT customer_id, password_hash, NOT $1 AND cardinality(tan_channels) > 0 FROM customers WHERE username=$2 AND deleted_at IS NULL ", - ) - .bind(is2fa) - .bind(username) - .fetch_optional(&mut *tx) - .await? - else { - return Ok(PatchAuthResult::UnknownAccount); - }; + ) + .bind(is2fa) + .bind(username) + .fetch_optional(&mut *tx) + .await? + else { + return Ok(PatchAuthResult::UnknownAccount); + }; - let res = if let Some(old_pw) = old_pw - && !pw_crypto.checkpw(old_pw, &currenc_pwh).unwrap().matches - { - PatchAuthResult::OldPasswordMismatch - } else if tan_required { - PatchAuthResult::TanRequired - } else { - let new_pwh = pw_crypto.hashpw(new_pw); - sqlx::query( + let res = if let Some(old_pw) = old_pw + && !pw_crypto.checkpw(old_pw, &currenc_pwh).unwrap().matches + { + PatchAuthResult::OldPasswordMismatch + } else if tan_required { + PatchAuthResult::TanRequired + } else { + let new_pwh = pw_crypto.hashpw(new_pw); + sqlx::query( "UPDATE customers SET password_hash=$1, token_creation_counter=0 WHERE customer_id=$2", ) .bind(new_pwh) .bind(customer_id) .execute(&mut *tx) .await?; - PatchAuthResult::Success - }; + PatchAuthResult::Success + }; - tx.commit().await?; - Ok(res) + tx.commit().await?; + Ok(res) + }) } /** Result status of customer account password check */ @@ -561,8 +572,9 @@ pub async fn check_password( pw: &str, ) -> sqlx::Result<CheckPasswordResult> { // Get user current password hash - let Some((info, pwh, counter)): Option<(_, CompactString, _)> = sqlx::query( - " + let Some((info, ref pwh, counter)): Option<(_, CompactString, _)> = serialized!( + sqlx::query( + " SELECT password_hash, token_creation_counter, @@ -577,26 +589,26 @@ pub async fn check_password( JOIN customers ON customer_id=owning_customer_id WHERE username=$1 AND deleted_at IS NULL ", - ) - .bind(username) - .try_map(|r: PgRow| { - let info = BankInfo { - username: username.into(), - payto: sql_bank_payto(&r, ctx, "internal_payto", "name")?, - bank_account_id: r.try_get_u64("bank_account_id")?, - is_exchange: r.try_get("is_taler_exchange")?, - phone: r.try_get("phone")?, - email: r.try_get("email")?, - channels: r.try_get("tan_channels")?, - }; - Ok(( - info, - r.try_get("password_hash")?, - r.try_get_u16("token_creation_counter")?, - )) - }) - .fetch_optional(db) - .await? + ) + .bind(username) + .try_map(|r: PgRow| { + let info = BankInfo { + username: username.into(), + payto: sql_bank_payto(&r, ctx, "internal_payto", "name")?, + bank_account_id: r.try_get_u64("bank_account_id")?, + is_exchange: r.try_get("is_taler_exchange")?, + phone: r.try_get("phone")?, + email: r.try_get("email")?, + channels: r.try_get("tan_channels")?, + }; + Ok(( + info, + r.try_get("password_hash")?, + r.try_get_u16("token_creation_counter")?, + )) + }) + .fetch_optional(db) + )? else { return Ok(CheckPasswordResult::UnknownAccount); }; @@ -607,20 +619,23 @@ pub async fn check_password( } // Check password - let check = pw_crypto.checkpw(pw, &pwh).unwrap(); // TODO handle this + let check = pw_crypto.checkpw(pw, pwh).unwrap(); // TODO handle this if !check.matches { return Ok(CheckPasswordResult::PasswordMismatch); } // Rehash if outdated if check.outdated { - let new = pw_crypto.hashpw(pw); - sqlx::query("UPDATE customers SET password_hash=$1 where username=$2 AND password_hash=$3") + let new = &pw_crypto.hashpw(pw); + serialized!( + sqlx::query( + "UPDATE customers SET password_hash=$1 where username=$2 AND password_hash=$3" + ) .bind(new) .bind(username) .bind(pwh) .execute(db) - .await?; + )?; } Ok(CheckPasswordResult::Success(info)) } @@ -662,8 +677,9 @@ pub async fn bank_info( ctx: &PaytoCtx, username: &str, ) -> sqlx::Result<Option<BankInfo>> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT bank_account_id, internal_payto, @@ -676,21 +692,21 @@ pub async fn bank_info( JOIN customers ON customer_id=owning_customer_id WHERE username=$1 ", - ) - .bind(username) - .try_map(|r: PgRow| { - Ok(BankInfo { - username: username.into(), - payto: sql_bank_payto(&r, ctx, "internal_payto", "name")?, - bank_account_id: r.try_get_u64("bank_account_id")?, - is_exchange: r.try_get("is_taler_exchange")?, - phone: r.try_get("phone")?, - email: r.try_get("email")?, - channels: r.try_get("tan_channels")?, + ) + .bind(username) + .try_map(|r: PgRow| { + Ok(BankInfo { + username: username.into(), + payto: sql_bank_payto(&r, ctx, "internal_payto", "name")?, + bank_account_id: r.try_get_u64("bank_account_id")?, + is_exchange: r.try_get("is_taler_exchange")?, + phone: r.try_get("phone")?, + email: r.try_get("email")?, + channels: r.try_get("tan_channels")?, + }) }) - }) - .fetch_optional(db) - .await + .fetch_optional(db) + ) } /** Check bank info of account [payto] */ @@ -711,8 +727,9 @@ pub async fn by_username( fiat: Option<&Currency>, username: &str, ) -> sqlx::Result<Option<AccountData>> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT customers.name, email, @@ -746,44 +763,44 @@ pub async fn by_username( CROSS JOIN LATERAL get_conversion_class_rate(conversion_rate_class_id) WHERE username=$2 ", - ) - .bind(MAX_TOKEN_CREATION_ATTEMPTS as i16) - .bind(username) - .try_map(|r: PgRow| { - let status: AccountStatus = r.try_get("status")?; - let channels: Vec<TanChannel> = r.try_get("tan_channels")?; - let is_exchange: bool = r.try_get("is_taler_exchange")?; - - Ok(AccountData { - payto_uri: sql_bank_payto(&r, ctx, "internal_payto", "name")?, - balance: Balance { - amount: r.try_get_amount("balance", regional)?, - credit_debit_indicator: if r.try_get("has_debt")? { - CreditDebitInfo::debit - } else { - CreditDebitInfo::credit + ) + .bind(MAX_TOKEN_CREATION_ATTEMPTS as i16) + .bind(username) + .try_map(|r: PgRow| { + let status: AccountStatus = r.try_get("status")?; + let channels: Vec<TanChannel> = r.try_get("tan_channels")?; + let is_exchange: bool = r.try_get("is_taler_exchange")?; + + Ok(AccountData { + payto_uri: sql_bank_payto(&r, ctx, "internal_payto", "name")?, + balance: Balance { + amount: r.try_get_amount("balance", regional)?, + credit_debit_indicator: if r.try_get("has_debt")? { + CreditDebitInfo::debit + } else { + CreditDebitInfo::credit + }, }, - }, - debit_threshold: r.try_get_amount("max_debt", regional)?, - contact_data: ChallengeContactData { - email: r.try_get("email")?, - phone: r.try_get("phone")?, - }, - cashout_payto_uri: sql_opt_iban_payto(&r, "cashout_payto", "name")? - .map(|it| it.as_uri()), - tan_channel: channels.first().cloned(), - tan_channels: channels, - is_public: r.try_get("is_public")?, - is_taler_exchange: is_exchange, - is_locked: status == AccountStatus::locked, - status, - conversion_rate_class_id: r.try_get_opt_u64("conversion_rate_class_id")?, - conversion_rate: user_rate(&r, regional, fiat, username, is_exchange)?, - name: r.try_get("name")?, + debit_threshold: r.try_get_amount("max_debt", regional)?, + contact_data: ChallengeContactData { + email: r.try_get("email")?, + phone: r.try_get("phone")?, + }, + cashout_payto_uri: sql_opt_iban_payto(&r, "cashout_payto", "name")? + .map(|it| it.as_uri()), + tan_channel: channels.first().cloned(), + tan_channels: channels, + is_public: r.try_get("is_public")?, + is_taler_exchange: is_exchange, + is_locked: status == AccountStatus::locked, + status, + conversion_rate_class_id: r.try_get_opt_u64("conversion_rate_class_id")?, + conversion_rate: user_rate(&r, regional, fiat, username, is_exchange)?, + name: r.try_get("name")?, + }) }) - }) - .fetch_optional(db) - .await + .fetch_optional(db) + ) } /** Get a page of all public accounts */ diff --git a/crates/libeufin-bank/src/db/exchange.rs b/crates/libeufin-bank/src/db/exchange.rs @@ -256,7 +256,7 @@ pub async fn add_incoming( .bind(debtor.canonical()) .bind(username) .bind_timestamp(timestamp) - .bind(metadata.ty()) + .bind(metadata.ty().as_ref()) .try_map(|r: PgRow| { Ok(if r.try_get_flag("out_creditor_not_found")? { AddIncomingResult::UnknownExchange @@ -321,7 +321,7 @@ pub async fn incoming_history( query }, |r| { - Ok(match r.try_get("type")? { + Ok(match r.try_get_parse("type")? { IncomingType::reserve => IncomingBankTransaction::Reserve { row_id: r.try_get_u64("bank_transaction_id")?, date: r.try_get_timestamp("transaction_date")?.into(), diff --git a/crates/libeufin-bank/src/db/prepared.rs b/crates/libeufin-bank/src/db/prepared.rs @@ -61,7 +61,7 @@ pub async fn register( " ) .bind(username) - .bind(ty) + .bind(ty.as_ref()) .bind(account_pub) .bind(auth_pub) .bind(auth_sig) diff --git a/crates/libeufin-bank/src/db/tan.rs b/crates/libeufin-bank/src/db/tan.rs @@ -24,7 +24,10 @@ use std::time::Duration; use compact_str::CompactString; use jiff::Timestamp; use sqlx::{PgPool, Row, postgres::PgRow}; -use taler_api::db::{BindHelper, TypeHelper}; +use taler_api::{ + db::{BindHelper, TypeHelper}, + serialized, +}; use taler_common::types::base32::Base32; use uuid::Uuid; @@ -44,8 +47,9 @@ pub async fn new( channel: TanChannel, info: &str, ) -> sqlx::Result<Uuid> { - sqlx::query_scalar( - " + serialized!( + sqlx::query_scalar( + " INSERT INTO tan_challenges ( hbody, salt, @@ -72,19 +76,19 @@ pub async fn new( gen_random_uuid() ) RETURNING uuid ", + ) + .bind(hash) + .bind(salt) + .bind(op) + .bind(code) + .bind_timestamp(now) + .bind_timestamp(&(*now + validity)) + .bind(retry_counter as i16) + .bind(username) + .bind(channel) + .bind(info) + .fetch_one(db) ) - .bind(hash) - .bind(salt) - .bind(op) - .bind(code) - .bind_timestamp(now) - .bind_timestamp(&(*now + validity)) - .bind(retry_counter as i16) - .bind(username) - .bind(channel) - .bind(info) - .fetch_one(db) - .await } /** Result of TAN challenge transmission */ @@ -112,7 +116,8 @@ pub async fn send( now: &Timestamp, max_active: u16, ) -> sqlx::Result<SendResult> { - Ok(sqlx::query(" + Ok(serialized!( + sqlx::query(" SELECT (confirmation_date IS NOT NULL) as solved ,retransmission_date @@ -160,16 +165,18 @@ pub async fn send( } } )) - .fetch_optional(db).await?.unwrap_or(SendResult::NotFound)) + .fetch_optional(db) + )?.unwrap_or(SendResult::NotFound)) } /** Mark TAN challenge transmission */ pub async fn mark_sent(db: &PgPool, uuid: &Uuid, retransmission: &Timestamp) -> sqlx::Result<()> { - sqlx::query("UPDATE tan_challenges SET retransmission_date = $1 WHERE uuid = $2") - .bind_timestamp(retransmission) - .bind(uuid) - .execute(db) - .await?; + serialized!( + sqlx::query("UPDATE tan_challenges SET retransmission_date = $1 WHERE uuid = $2") + .bind_timestamp(retransmission) + .bind(uuid) + .execute(db) + )?; Ok(()) } @@ -193,36 +200,37 @@ pub async fn solve( code: &str, timestamp: &Timestamp, ) -> sqlx::Result<SolveResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_ok, out_no_op, out_no_retry, out_expired, out_op, out_channel, out_info FROM tan_challenge_try($1,$2,$3) ", - ) - .bind(uuid) - .bind(code) - .bind_timestamp(timestamp) - .try_map(|r: PgRow| { - Ok(if r.try_get_flag("out_ok")? { - SolveResult::Success { - op: r.try_get("out_op")?, - channel: r.try_get("out_channel")?, - info: r.try_get("out_info")?, - } - } else if r.try_get_flag("out_no_op")? { - SolveResult::NotFound - } else if r.try_get_flag("out_no_retry")? { - SolveResult::NoRetry - } else if r.try_get_flag("out_expired")? { - SolveResult::Expired - } else { - SolveResult::BadCode + ) + .bind(uuid) + .bind(code) + .bind_timestamp(timestamp) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_ok")? { + SolveResult::Success { + op: r.try_get("out_op")?, + channel: r.try_get("out_channel")?, + info: r.try_get("out_info")?, + } + } else if r.try_get_flag("out_no_op")? { + SolveResult::NotFound + } else if r.try_get_flag("out_no_retry")? { + SolveResult::NoRetry + } else if r.try_get_flag("out_expired")? { + SolveResult::Expired + } else { + SolveResult::BadCode + }) }) - }) - .fetch_one(db) - .await + .fetch_one(db) + ) } #[derive(Debug)] @@ -237,25 +245,26 @@ pub struct SolvedChallenge { } pub async fn challenge(db: &PgPool, uuids: &[Uuid]) -> sqlx::Result<Vec<SolvedChallenge>> { - sqlx::query( - " - SELECT uuid, salt, hbody, tan_channel, tan_info, op, (confirmation_date IS NOT NULL) as confirmed - FROM tan_challenges - WHERE uuid = ANY($1) - ", - ) - .bind(uuids) - .try_map(|r: PgRow| { - Ok(SolvedChallenge { - id: r.try_get("uuid")?, - salt: r.try_get("salt")?, - hash: r.try_get("hbody")?, - channel: r.try_get("tan_channel")?, - info: r.try_get("tan_info")?, - confirmed: r.try_get("confirmed")?, - op: r.try_get("op")?, + serialized!( + sqlx::query( + " + SELECT uuid, salt, hbody, tan_channel, tan_info, op, (confirmation_date IS NOT NULL) as confirmed + FROM tan_challenges + WHERE uuid = ANY($1) + ", + ) + .bind(uuids) + .try_map(|r: PgRow| { + Ok(SolvedChallenge { + id: r.try_get("uuid")?, + salt: r.try_get("salt")?, + hash: r.try_get("hbody")?, + channel: r.try_get("tan_channel")?, + info: r.try_get("tan_info")?, + confirmed: r.try_get("confirmed")?, + op: r.try_get("op")?, + }) }) - }) - .fetch_all(db) - .await + .fetch_all(db) + ) } diff --git a/crates/libeufin-bank/src/db/tx.rs b/crates/libeufin-bank/src/db/tx.rs @@ -59,7 +59,7 @@ pub async fn create( amount: Amount, timestamp: &Timestamp, is2fa: bool, - request_uid: Option<ShortHashCode>, + request_uid: &Option<ShortHashCode>, wire_transfer_fees: Amount, min_amount: Amount, max_amount: Amount, @@ -78,8 +78,9 @@ pub async fn create( Ok(s) => (Some(s.ty()), Some(s.key()), None), Err(e) => (None, None, Some(e)), }; - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_creditor_not_found ,out_debtor_not_found @@ -98,45 +99,45 @@ pub async fn create( ,out_idempotent FROM bank_transaction($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11::taler_incoming_type,$12,$13) ", - ) - .bind(credit_account_payto.canonical()) - .bind(debit_account_username) - .bind(subject) - .bind(amount) - .bind_timestamp(timestamp) - .bind(is2fa) - .bind(request_uid) - .bind(wire_transfer_fees) - .bind(min_amount) - .bind(max_amount) - .bind(ty) - .bind(metadata) - .bind(cause) - .try_map(|r: PgRow| { - Ok(if r.try_get_flag("out_creditor_not_found")? { - TxResult::UnknownCreditor - } else if r.try_get_flag("out_debtor_not_found")? { - TxResult::UnknownDebtor - } else if r.try_get_flag("out_same_account")? { - TxResult::BothPartySame - } else if r.try_get_flag("out_balance_insufficient")? { - TxResult::BalanceInsufficient - } else if r.try_get_flag("out_bad_amount")? { - TxResult::BadAmount - } else if r.try_get_flag("out_creditor_admin")? { - TxResult::AdminCreditor - } else if r.try_get_flag("out_request_uid_reuse")? { - TxResult::RequestUidReuse - } else if r.try_get_flag("out_idempotent")? { - TxResult::Success(r.try_get_u64("out_debit_row_id")?) - } else if r.try_get_flag("out_tan_required")? { - TxResult::TanRequired - } else { - TxResult::Success(r.try_get_u64("out_debit_row_id")?) + ) + .bind(credit_account_payto.canonical()) + .bind(debit_account_username) + .bind(subject) + .bind(amount) + .bind_timestamp(timestamp) + .bind(is2fa) + .bind(request_uid) + .bind(wire_transfer_fees) + .bind(min_amount) + .bind(max_amount) + .bind(ty.as_ref().map(|it| it.as_ref())) + .bind(metadata) + .bind(cause) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_creditor_not_found")? { + TxResult::UnknownCreditor + } else if r.try_get_flag("out_debtor_not_found")? { + TxResult::UnknownDebtor + } else if r.try_get_flag("out_same_account")? { + TxResult::BothPartySame + } else if r.try_get_flag("out_balance_insufficient")? { + TxResult::BalanceInsufficient + } else if r.try_get_flag("out_bad_amount")? { + TxResult::BadAmount + } else if r.try_get_flag("out_creditor_admin")? { + TxResult::AdminCreditor + } else if r.try_get_flag("out_request_uid_reuse")? { + TxResult::RequestUidReuse + } else if r.try_get_flag("out_idempotent")? { + TxResult::Success(r.try_get_u64("out_debit_row_id")?) + } else if r.try_get_flag("out_tan_required")? { + TxResult::TanRequired + } else { + TxResult::Success(r.try_get_u64("out_debit_row_id")?) + }) }) - }) - .fetch_one(db) - .await + .fetch_one(db) + ) } // /** Get transaction [rowId] owned by [username] */ diff --git a/crates/libeufin-ebics/src/db.rs b/crates/libeufin-ebics/src/db.rs @@ -55,12 +55,12 @@ pub mod test { let ids = ["first", "second", "third"]; for id in ids { - ebics_register(&db, id).await.unwrap(); + 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!(Some(id), ebics_first(db).await.unwrap().as_deref()); + ebics_remove(db, id).await.unwrap(); } - assert_eq!(ebics_first(&db).await.unwrap(), None); + assert_eq!(ebics_first(db).await.unwrap(), None); } } diff --git a/crates/libeufin-ebics/src/iso20022/pain001.rs b/crates/libeufin-ebics/src/iso20022/pain001.rs @@ -202,7 +202,7 @@ mod test { ); let msg = Pain001Msg { - msg_id: "MESSAGE_ID".into(), + msg_id: "MESSAGE_ID", timestamp: &date_to_timestamp("2024-09-09"), debtor: &FullIbanPayto::new( BankID { @@ -222,11 +222,11 @@ mod test { Pain001Tx { creditor: creditor.clone(), amount: amount("CHF:5.11"), - subject: "Test 5.11".into(), + subject: "Test 5.11", e2e_id: "TX_SECOND", }, Pain001Tx { - creditor: creditor, + creditor, amount: amount("CHF:0.21"), subject: "Test 0.21", e2e_id: "TX_THIRD", diff --git a/crates/libeufin-ebics/src/ws.rs b/crates/libeufin-ebics/src/ws.rs @@ -378,7 +378,7 @@ mod test { #[tokio::test] pub async fn params() { let path = "/tmp/libeufin_nexus_wss_test.sock"; - std::fs::remove_file(&path).ok(); + std::fs::remove_file(path).ok(); let server = axum::Router::new() .route( "/", @@ -405,7 +405,7 @@ mod test { .serve( taler_api::Serve::Unix { path: path.into(), - permission: Permissions::from_mode(660), + permission: Permissions::from_mode(0o660), }, None, ); diff --git a/crates/libeufin-nexus/src/bench.rs b/crates/libeufin-nexus/src/bench.rs @@ -187,7 +187,7 @@ mod test { b.measure("wg_transfer", async |_| { server .post("/taler-wire-gateway/transfer") - .json(&json!({ + .json(json!({ "request_uid": HashCode::rand(), "amount": "KUDOS:0.0001", "exchange_base_url": "http://exchange.example.com/", @@ -200,7 +200,7 @@ mod test { .await; b.measure("wg_transfer_get", async |i| { server - .get(&format!("/taler-wire-gateway/transfers/{}", i + 1)) + .get(format!("/taler-wire-gateway/transfers/{}", i + 1)) .await .assert_ok() }) @@ -222,7 +222,7 @@ mod test { b.measure("wg_add", async |_| { server .post("/taler-wire-gateway/admin/add-incoming") - .json(&json!({ + .json(json!({ "amount": "KUDOS:0.0001", "reserve_pub": EddsaPublicKey::rand(), "debit_account": &*ACCOUNT @@ -252,13 +252,13 @@ mod test { server .post("/taler-prepared-transfer/registration") - .json(&json!({ + .json(json!({ "credit_amount": "KUDOS:55", "type": "reserve", "alg": "EdDSA", "account_pub": key, "authorization_pub": key, - "authorization_sig": eddsa_sign(&pair, key.as_ref()), + "authorization_sig": eddsa_sign(pair, key.as_ref()), "recurrent":false })) .await @@ -271,7 +271,7 @@ mod test { let req = json!({ "timestamp": &now, "authorization_pub": key, - "authorization_sig": eddsa_sign(&pair, now.as_ref()), + "authorization_sig": eddsa_sign(pair, now.as_ref()), }); server .post("/taler-prepared-transfer/unregistration") diff --git a/crates/libeufin-nexus/src/db.rs b/crates/libeufin-nexus/src/db.rs @@ -186,7 +186,7 @@ pub mod test { } else if r.try_get_flag(2)? { Status::Incomplete } else { - match r.try_get(3)? { + match r.try_get_opt_parse(3)? { None => Status::Simple, Some(IncomingType::reserve) => Status::Reserve(r.try_get(4)?), Some(IncomingType::kyc) => Status::Kyc(r.try_get(4)?), diff --git a/crates/libeufin-nexus/src/db/exchange.rs b/crates/libeufin-nexus/src/db/exchange.rs @@ -118,7 +118,7 @@ pub async fn incoming_history( let credit_fee = r .try_get_opt_amount("credit_fee", currency)? .filter(|it| it.is_zero()); - Ok(match r.try_get("type")? { + Ok(match r.try_get_parse("type")? { IncomingType::reserve => IncomingBankTransaction::Reserve { row_id: r.try_get_u64("incoming_transaction_id")?, amount: r.try_get_amount("amount", currency)?, diff --git a/crates/libeufin-nexus/src/db/initiated.rs b/crates/libeufin-nexus/src/db/initiated.rs @@ -20,7 +20,10 @@ use const_format::formatcp; use jiff::Timestamp; use libeufin_ebics::iso20022::model::{OutId, OutTx}; use sqlx::{PgPool, Row as _, postgres::PgRow}; -use taler_api::db::{BindHelper as _, TypeHelper as _}; +use taler_api::{ + db::{BindHelper as _, TypeHelper as _}, + serialized, +}; use taler_common::types::{ amount::{Amount, Currency}, payto::PaytoURI, @@ -47,8 +50,9 @@ pub async fn initiate( initiation_time: &Timestamp, e2e_id: &str, ) -> sqlx::Result<PaymentInitiationResult> { - let res = sqlx::query( - " + let res = serialized!( + sqlx::query( + " INSERT INTO initiated_outgoing_transactions ( amount, subject, @@ -58,15 +62,15 @@ pub async fn initiate( ) VALUES ($1,$2,$3,$4,$5) RETURNING initiated_outgoing_transaction_id ", - ) - .bind(amount) - .bind(subject) - .bind(creditor.as_ref().as_str()) - .bind_timestamp(initiation_time) - .bind(e2e_id) - .try_map(|r: PgRow| Ok(PaymentInitiationResult::Success(r.try_get_u64(0)?))) - .fetch_one(pool) - .await; + ) + .bind(amount) + .bind(subject) + .bind(creditor.as_ref().as_str()) + .bind_timestamp(initiation_time) + .bind(e2e_id) + .try_map(|r: PgRow| Ok(PaymentInitiationResult::Success(r.try_get_u64(0)?))) + .fetch_one(pool) + ); if let Err(e) = &res && let Some(db_err) = e.as_database_error() && db_err.code() == Some(std::borrow::Cow::Borrowed("23505")) @@ -84,20 +88,22 @@ pub async fn batch_initiated( ebics_id: &str, require_ack: bool, ) -> sqlx::Result<()> { - sqlx::query("SELECT batch_outgoing_transactions($1, $2, $3)") - .bind_timestamp(timestamp) - .bind(ebics_id) - .bind(require_ack) - .execute(pool) - .await?; + serialized!( + sqlx::query("SELECT batch_outgoing_transactions($1, $2, $3)") + .bind_timestamp(timestamp) + .bind(ebics_id) + .bind(require_ack) + .execute(pool) + )?; Ok(()) } pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<()> { - sqlx::query("UPDATE initiated_outgoing_transactions SET awaiting_ack=false WHERE initiated_outgoing_transaction_id=$1") + serialized!( + sqlx::query("UPDATE initiated_outgoing_transactions SET awaiting_ack=false WHERE initiated_outgoing_transaction_id=$1") .bind(id as i64) .execute(db) - .await?; + )?; Ok(()) } @@ -109,37 +115,38 @@ pub async fn initiated_submittable( SELECT initiated_outgoing_batch_id, message_id, creation_date, sum FROM initiated_outgoing_batches "; - let mut tx = db.begin().await?; - // We want to maximize the number of successfully submitted batches in the event - // of a malformed transaction or a persistent error classified as transient. We send - // the unsubmitted batches first, starting with the oldest by creation time. - // This is the happy path, giving every batch a chance while being fair on the - // basis of creation date. - // Then we retry the failed batches, starting with the oldest by submission time. - // This the bad path retrying each failed batch applying a rotation based on - // resubmission time. - let mut batches = sqlx::query(formatcp!( - " + serialized!(async { + let mut tx = db.begin().await?; + // We want to maximize the number of successfully submitted batches in the event + // of a malformed transaction or a persistent error classified as transient. We send + // the unsubmitted batches first, starting with the oldest by creation time. + // This is the happy path, giving every batch a chance while being fair on the + // basis of creation date. + // Then we retry the failed batches, starting with the oldest by submission time. + // This the bad path retrying each failed batch applying a rotation based on + // resubmission time. + let mut batches = sqlx::query(formatcp!( + " ({SELECT_PART} WHERE status='unsubmitted' ORDER BY creation_date ASC) UNION ALL ({SELECT_PART} WHERE status='transient_failure' ORDER BY submission_date) " - )) - .try_map(|r: PgRow| { - Ok(PaymentBatch { - id: r.try_get_u64("initiated_outgoing_batch_id")?, - msg_id: r.try_get("message_id")?, - creation_date: r.try_get_timestamp("creation_date")?, - sum: r.try_get_amount("sum", currency)?, - payments: Vec::new(), + )) + .try_map(|r: PgRow| { + Ok(PaymentBatch { + id: r.try_get_u64("initiated_outgoing_batch_id")?, + msg_id: r.try_get("message_id")?, + creation_date: r.try_get_timestamp("creation_date")?, + sum: r.try_get_amount("sum", currency)?, + payments: Vec::new(), + }) }) - }) - .fetch_all(&mut *tx) - .await?; - let mut batch_map: BTreeMap<_, _> = batches.iter_mut().map(|it| (it.id, it)).collect(); - // Then load transactions - sqlx::query( - " + .fetch_all(&mut *tx) + .await?; + let mut batch_map: BTreeMap<_, _> = batches.iter_mut().map(|it| (it.id, it)).collect(); + // Then load transactions + sqlx::query( + " SELECT initiated_outgoing_transaction_id ,amount @@ -152,24 +159,25 @@ pub async fn initiated_submittable( JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id) WHERE initiated_outgoing_batches.status IN ('unsubmitted', 'transient_failure') ", - ) - .try_map(|r: PgRow| { - let payment = Initiated { - id: r.try_get_u64("initiated_outgoing_transaction_id")?, - amount: r.try_get_amount("amount", currency)?, - creditor: r.try_get_parse("credit_payto")?, - subject: r.try_get("subject")?, - initiation_time: r.try_get_timestamp("initiation_time")?, - e2e_id: r.try_get("end_to_end_id")?, - }; - let batch_id = r.try_get_u64("initiated_outgoing_batch_id")?; - batch_map.get_mut(&batch_id).unwrap().payments.push(payment); - Ok(()) + ) + .try_map(|r: PgRow| { + let payment = Initiated { + id: r.try_get_u64("initiated_outgoing_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + creditor: r.try_get_parse("credit_payto")?, + subject: r.try_get("subject")?, + initiation_time: r.try_get_timestamp("initiation_time")?, + e2e_id: r.try_get("end_to_end_id")?, + }; + let batch_id = r.try_get_u64("initiated_outgoing_batch_id")?; + batch_map.get_mut(&batch_id).unwrap().payments.push(payment); + Ok(()) + }) + .fetch_all(&mut *tx) + .await?; + tx.commit().await?; + Ok(batches) }) - .fetch_all(&mut *tx) - .await?; - tx.commit().await?; - Ok(batches) } pub async fn unsettled_tx_in_batch( @@ -178,9 +186,10 @@ pub async fn unsettled_tx_in_batch( msg_id: &str, execution_time: &Timestamp, ) -> sqlx::Result<Vec<OutTx>> { - sqlx::query(formatcp!( - " - SELECT + serialized!( + sqlx::query(formatcp!( + " + SELECT end_to_end_id, amount, subject, @@ -190,24 +199,24 @@ pub async fn unsettled_tx_in_batch( WHERE message_id = $1 AND initiated_outgoing_transactions.{UNSETTLED} " - )) - .bind(msg_id) - .try_map(|r: PgRow| { - Ok(OutTx { - id: OutId { - msg_id: Some(msg_id.into()), - e2e_id: r.try_get("end_to_end_id")?, - sref: None, - }, - amount: r.try_get_amount("amount", currency)?, - debit_fee: Amount::zero(currency), - subject: r.try_get("subject")?, - execution_time: *execution_time, - creditor: r.try_get_opt_payto("credit_payto")?, + )) + .bind(msg_id) + .try_map(|r: PgRow| { + Ok(OutTx { + id: OutId { + msg_id: Some(msg_id.into()), + e2e_id: r.try_get("end_to_end_id")?, + sref: None, + }, + amount: r.try_get_amount("amount", currency)?, + debit_fee: Amount::zero(currency), + subject: r.try_get("subject")?, + execution_time: *execution_time, + creditor: r.try_get_opt_payto("credit_payto")?, + }) }) - }) - .fetch_all(db) - .await + .fetch_all(db) + ) } /** Register submission success of order [orderId] for batch [id] at [timestamp] */ @@ -217,11 +226,12 @@ pub async fn batch_sub_success( timestamp: &Timestamp, order_id: &str, ) -> sqlx::Result<()> { - let mut tx = db.begin().await?; - // Update batch status - let updated = sqlx::query( - " - UPDATE initiated_outgoing_batches + serialized!(async { + let mut tx = db.begin().await?; + // Update batch status + let updated = sqlx::query( + " + UPDATE initiated_outgoing_batches SET status = 'pending' ,submission_date = $1 ,status_msg = NULL @@ -229,26 +239,27 @@ pub async fn batch_sub_success( ,submission_counter = submission_counter + 1 WHERE initiated_outgoing_batch_id = $3 AND order_id IS NULL ", - ) - .bind_timestamp(timestamp) - .bind(order_id) - .bind(batch_id as i64) - .execute(&mut *tx) - .await?; - if updated.rows_affected() > 0 { - // Update unsettled batch's transaction status - sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_transactions - SET status = 'pending', status_msg = NULL - WHERE initiated_outgoing_batch_id = $1 AND {UNSETTLED} - " - )) + ) + .bind_timestamp(timestamp) + .bind(order_id) .bind(batch_id as i64) .execute(&mut *tx) .await?; - } - tx.commit().await + if updated.rows_affected() > 0 { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions + SET status = 'pending', status_msg = NULL + WHERE initiated_outgoing_batch_id = $1 AND {UNSETTLED} + " + )) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await + }) } /** Register submission failure with [msg] for batch [id] at [timestamp]*/ @@ -259,10 +270,11 @@ pub async fn batch_sub_failure( msg: &str, ) -> sqlx::Result<()> { let permanent = false; - let mut tx = db.begin().await?; - // Update batch status - sqlx::query( - " + serialized!(async { + let mut tx = db.begin().await?; + // Update batch status + sqlx::query( + " UPDATE initiated_outgoing_batches SET status = $1 ,submission_date = $2 @@ -270,102 +282,107 @@ pub async fn batch_sub_failure( ,submission_counter = submission_counter + 1 WHERE initiated_outgoing_batch_id = $4 ", - ) - .bind(if permanent { - SubmissionState::permanent_failure - } else { - SubmissionState::transient_failure - }) - .bind_timestamp(timestamp) - .bind(msg) - .bind(batch_id as i64) - .execute(&mut *tx) - .await?; - // Update unsettled batch's transaction status - sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_transactions + ) + .bind(if permanent { + SubmissionState::permanent_failure + } else { + SubmissionState::transient_failure + }) + .bind_timestamp(timestamp) + .bind(msg) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions SET status = $1, status_msg = $2 WHERE initiated_outgoing_batch_id = $3 AND {UNSETTLED} " - )) - .bind(if permanent { - SubmissionState::permanent_failure - } else { - SubmissionState::transient_failure + )) + .bind(if permanent { + SubmissionState::permanent_failure + } else { + SubmissionState::transient_failure + }) + .bind(msg) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + tx.commit().await }) - .bind(msg) - .bind(batch_id as i64) - .execute(&mut *tx) - .await?; - tx.commit().await } /** Register order step [msg] for [orderId] */ pub async fn order_step(db: &PgPool, order_id: &str, msg: &str) -> sqlx::Result<()> { - let mut tx = db.begin().await?; - // Update batch status - let batch_id = sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_batches + serialized!(async { + let mut tx = db.begin().await?; + // Update batch status + let batch_id = sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_batches SET status = 'pending', status_msg = $1 WHERE order_id = $2 AND {PENDING} RETURNING initiated_outgoing_batch_id " - )) - .bind(msg) - .bind(order_id) - .try_map(|r: PgRow| r.try_get_u64(0)) - .fetch_optional(&mut *tx) - .await?; - if let Some(batch_id) = batch_id { - // Update unsettled batch's transaction status - sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_transactions - SET status = 'pending', status_msg = $1 - WHERE initiated_outgoing_batch_id = $2 AND {PENDING} - " )) .bind(msg) - .bind(batch_id as i64) - .execute(&mut *tx) + .bind(order_id) + .try_map(|r: PgRow| r.try_get_u64(0)) + .fetch_optional(&mut *tx) .await?; - } - tx.commit().await + if let Some(batch_id) = batch_id { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions + SET status = 'pending', status_msg = $1 + WHERE initiated_outgoing_batch_id = $2 AND {PENDING} + " + )) + .bind(msg) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await + }) } /** Register order success for [orderId] and return message_id if found */ pub async fn order_success(db: &PgPool, order_id: &str) -> sqlx::Result<Option<String>> { - let mut tx = db.begin().await?; - // Update batch status - let res = sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_batches + serialized!(async { + let mut tx = db.begin().await?; + // Update batch status + let res = sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_batches SET status = 'success' WHERE order_id = $1 RETURNING initiated_outgoing_batch_id, message_id " - )) - .bind(order_id) - .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?))) - .fetch_optional(&mut *tx) - .await?; - if let Some((batch_id, _)) = &res { - // Update unsettled batch's transaction status - sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_transactions + )) + .bind(order_id) + .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?))) + .fetch_optional(&mut *tx) + .await?; + if let Some((batch_id, _)) = &res { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions SET status = 'pending' WHERE initiated_outgoing_batch_id = $1 AND {UNSETTLED} " - )) - .bind(*batch_id as i64) - .execute(&mut *tx) - .await?; - } - tx.commit().await?; - Ok(res.map(|(_, msg_id)| msg_id)) + )) + .bind(*batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await?; + Ok(res.map(|(_, msg_id)| msg_id)) + }) } /** Register order failure for [orderId] and return message_id and previous status_msg if found */ @@ -373,35 +390,37 @@ pub async fn order_failure( db: &PgPool, order_id: &str, ) -> sqlx::Result<Option<(String, Option<String>)>> { - let mut tx = db.begin().await?; - // Update batch status - let res = sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_batches + serialized!(async { + let mut tx = db.begin().await?; + // Update batch status + let res = sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_batches SET status = 'permanent_failure' WHERE order_id = $1 RETURNING initiated_outgoing_batch_id, message_id, status_msg " - )) - .bind(order_id) - .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?, r.try_get(2)?))) - .fetch_optional(&mut *tx) - .await?; - if let Some((batch_id, _, _)) = &res { - // Update unsettled batch's transaction status - sqlx::query(formatcp!( - " - UPDATE initiated_outgoing_transactions + )) + .bind(order_id) + .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?, r.try_get(2)?))) + .fetch_optional(&mut *tx) + .await?; + if let Some((batch_id, _, _)) = &res { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions SET status = 'permanent_failure' WHERE initiated_outgoing_batch_id = $1 " - )) - .bind(*batch_id as i64) - .execute(&mut *tx) - .await?; - } - tx.commit().await?; - Ok(res.map(|(_, msg_id, status_msg)| (msg_id, status_msg))) + )) + .bind(*batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await?; + Ok(res.map(|(_, msg_id, status_msg)| (msg_id, status_msg))) + }) } /** Register payment status [state] with [msg] for batch [msgId] */ @@ -411,15 +430,16 @@ pub async fn batch_status_update( state: SubmissionState, msg: &str, ) -> sqlx::Result<bool> { - sqlx::query(formatcp!( - "SELECT out_ok FROM batch_status_update($1,$2,$3)" - )) - .bind(msg_id) - .bind(state) - .bind(msg) - .try_map(|r: PgRow| r.try_get(0)) - .fetch_one(db) - .await + serialized!( + sqlx::query(formatcp!( + "SELECT out_ok FROM batch_status_update($1,$2,$3)" + )) + .bind(msg_id) + .bind(state) + .bind(msg) + .try_map(|r: PgRow| r.try_get(0)) + .fetch_one(db) + ) } /** Register payment status [state] with [msg] for transaction [endToEndId] in batch [msgId] */ @@ -430,16 +450,17 @@ pub async fn tx_status_update( state: SubmissionState, msg: &str, ) -> sqlx::Result<bool> { - sqlx::query(formatcp!( - "SELECT out_ok FROM tx_status_update($1,$2,$3,$4)" - )) - .bind(end_to_end_id) - .bind(msg_id) - .bind(state) - .bind(msg) - .try_map(|r: PgRow| r.try_get(0)) - .fetch_one(db) - .await + serialized!( + sqlx::query(formatcp!( + "SELECT out_ok FROM tx_status_update($1,$2,$3,$4)" + )) + .bind(end_to_end_id) + .bind(msg_id) + .bind(state) + .bind(msg) + .try_map(|r: PgRow| r.try_get(0)) + .fetch_one(db) + ) } #[cfg(test)] @@ -649,7 +670,7 @@ mod test { sqlx::query( " SELECT (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_batches WHERE message_id != 'BATCH') - AND (SELECT bool_and(initiated_outgoing_transactions.status = 'unsubmitted') + AND (SELECT bool_and(initiated_outgoing_transactions.status = 'unsubmitted') FROM initiated_outgoing_transactions JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id) WHERE message_id != 'BATCH') " diff --git a/crates/libeufin-nexus/src/db/list.rs b/crates/libeufin-nexus/src/db/list.rs @@ -18,7 +18,7 @@ use compact_str::CompactString; use jiff::Timestamp; use libeufin_ebics::iso20022::model::{InId, OutId}; use sqlx::{PgPool, Row as _, postgres::PgRow}; -use taler_api::db::TypeHelper as _; +use taler_api::{db::TypeHelper as _, serialized}; use taler_common::{ api::{EddsaPublicKey, ShortHashCode}, types::amount::{Amount, Currency, Decimal}, @@ -126,44 +126,46 @@ pub async fn incoming( ORDER BY execution_time " }; - sqlx::query(query) - .try_map(|r: PgRow| { - let auth_pub: Option<EddsaPublicKey> = r.try_get("auth_pub")?; - let pending_pub: Option<EddsaPublicKey> = r.try_get("pending_pub")?; - let map = if let Some(auth_pub) = auth_pub { - format!(" mapped by {auth_pub}") - } else { - String::new() - }; - Ok(InMetadata { - id: InId { - uetr: r.try_get("uetr")?, - tx_id: r.try_get("tx_id")?, - sref: r.try_get("acct_svcr_ref")?, - }, - date: r.try_get_timestamp("execution_time")?, - amount: r.try_get_amount("amount", currency)?, - credit_fee: r.try_get("credit_fee")?, - subject: r.try_get("subject")?, - debtor: r.try_get("debit_payto")?, - bounced: r.try_get("bounced")?, - talerable: match r.try_get::<Option<CompactString>, _>("type")? { - None => pending_pub.map(|pending| format!("pending mapped by {pending}")), - Some(ty) => Some(format!( - "{ty} {}{map}", - r.try_get::<EddsaPublicKey, _>("metadata")? - )), - }, + serialized!( + sqlx::query(query) + .try_map(|r: PgRow| { + let auth_pub: Option<EddsaPublicKey> = r.try_get("auth_pub")?; + let pending_pub: Option<EddsaPublicKey> = r.try_get("pending_pub")?; + let map = if let Some(auth_pub) = auth_pub { + format!(" mapped by {auth_pub}") + } else { + String::new() + }; + Ok(InMetadata { + id: InId { + uetr: r.try_get("uetr")?, + tx_id: r.try_get("tx_id")?, + sref: r.try_get("acct_svcr_ref")?, + }, + date: r.try_get_timestamp("execution_time")?, + amount: r.try_get_amount("amount", currency)?, + credit_fee: r.try_get("credit_fee")?, + subject: r.try_get("subject")?, + debtor: r.try_get("debit_payto")?, + bounced: r.try_get("bounced")?, + talerable: match r.try_get::<Option<CompactString>, _>("type")? { + None => pending_pub.map(|pending| format!("pending mapped by {pending}")), + Some(ty) => Some(format!( + "{ty} {}{map}", + r.try_get::<EddsaPublicKey, _>("metadata")? + )), + }, + }) }) - }) - .fetch_all(db) - .await + .fetch_all(db) + ) } /** List outgoing transaction metadata for debugging */ pub async fn outgoing(db: &PgPool, currency: &Currency) -> sqlx::Result<Vec<OutMetadata>> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT amount ,subject @@ -177,30 +179,31 @@ pub async fn outgoing(db: &PgPool, currency: &Currency) -> sqlx::Result<Vec<OutM LEFT JOIN talerable_outgoing_transactions using (outgoing_transaction_id) ORDER BY execution_time ", - ) - .try_map(|r: PgRow| { - Ok(OutMetadata { - id: OutId { - msg_id: None, - e2e_id: r.try_get("end_to_end_id")?, - sref: r.try_get("acct_svcr_ref")?, - }, - date: r.try_get_timestamp("execution_time")?, - amount: r.try_get_amount("amount", currency)?, - subject: r.try_get("subject")?, - creditor: r.try_get("credit_payto")?, - wtid: r.try_get("wtid")?, - exchange_base_url: r.try_get("exchange_base_url")?, + ) + .try_map(|r: PgRow| { + Ok(OutMetadata { + id: OutId { + msg_id: None, + e2e_id: r.try_get("end_to_end_id")?, + sref: r.try_get("acct_svcr_ref")?, + }, + date: r.try_get_timestamp("execution_time")?, + amount: r.try_get_amount("amount", currency)?, + subject: r.try_get("subject")?, + creditor: r.try_get("credit_payto")?, + wtid: r.try_get("wtid")?, + exchange_base_url: r.try_get("exchange_base_url")?, + }) }) - }) - .fetch_all(db) - .await + .fetch_all(db) + ) } /** List initiated transaction metadata for debugging */ pub async fn initiated(db: &PgPool, currency: &Currency) -> sqlx::Result<Vec<InitMetadata>> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT amount ,subject @@ -217,30 +220,31 @@ pub async fn initiated(db: &PgPool, currency: &Currency) -> sqlx::Result<Vec<Ini LEFT JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id) ORDER BY initiation_time ", - ) - .try_map(|r: PgRow| { - Ok(InitMetadata { - date: r.try_get_timestamp("initiation_time")?, - amount: r.try_get_amount("amount", currency)?, - subject: r.try_get("subject")?, - creditor: r.try_get("credit_payto")?, - id: r.try_get("end_to_end_id")?, - batch: r.try_get("message_id")?, - batch_order: r.try_get("order_id")?, - status: r.try_get("status")?, - msg: r.try_get("status_msg")?, - submission_time: r.try_get_opt_timestamp("submission_date")?, - submission_counter: r.try_get_opt_u32("submission_counter")?.unwrap_or_default(), + ) + .try_map(|r: PgRow| { + Ok(InitMetadata { + date: r.try_get_timestamp("initiation_time")?, + amount: r.try_get_amount("amount", currency)?, + subject: r.try_get("subject")?, + creditor: r.try_get("credit_payto")?, + id: r.try_get("end_to_end_id")?, + batch: r.try_get("message_id")?, + batch_order: r.try_get("order_id")?, + status: r.try_get("status")?, + msg: r.try_get("status_msg")?, + submission_time: r.try_get_opt_timestamp("submission_date")?, + submission_counter: r.try_get_opt_u32("submission_counter")?.unwrap_or_default(), + }) }) - }) - .fetch_all(db) - .await + .fetch_all(db) + ) } /** List initiated transaction metadata pending acknowledgment for debugging */ pub async fn initiated_ack(db: &PgPool, currency: &Currency) -> sqlx::Result<Vec<InitMetadataAck>> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT amount ,subject @@ -252,17 +256,17 @@ pub async fn initiated_ack(db: &PgPool, currency: &Currency) -> sqlx::Result<Vec WHERE initiated_outgoing_batch_id IS NULL AND NOT awaiting_ack ORDER BY initiation_time ", - ) - .try_map(|r: PgRow| { - Ok(InitMetadataAck { - date: r.try_get_timestamp("initiation_time")?, - amount: r.try_get_amount("amount", currency)?, - subject: r.try_get("subject")?, - creditor: r.try_get("credit_payto")?, - id: r.try_get("end_to_end_id")?, - db_id: r.try_get_u64("initiated_outgoing_transaction_id")?, + ) + .try_map(|r: PgRow| { + Ok(InitMetadataAck { + date: r.try_get_timestamp("initiation_time")?, + amount: r.try_get_amount("amount", currency)?, + subject: r.try_get("subject")?, + creditor: r.try_get("credit_payto")?, + id: r.try_get("end_to_end_id")?, + db_id: r.try_get_u64("initiated_outgoing_transaction_id")?, + }) }) - }) - .fetch_all(db) - .await + .fetch_all(db) + ) } diff --git a/crates/libeufin-nexus/src/db/payment.rs b/crates/libeufin-nexus/src/db/payment.rs @@ -20,6 +20,7 @@ use libeufin_ebics::iso20022::model::{InTx, OutTx}; use sqlx::{PgPool, Row as _, postgres::PgRow}; use taler_api::{ db::{BindHelper as _, TypeHelper as _}, + serialized, subject::{IncomingSubject, OutgoingSubject}, }; use taler_common::types::amount::Amount; @@ -37,32 +38,33 @@ pub async fn register_out_tx( payment: &OutTx, subject: Option<&OutgoingSubject>, ) -> sqlx::Result<OutgoingRegistrationResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_tx_id, out_initiated, out_found FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) ", - ) - .bind(payment.amount) - .bind(payment.debit_fee) - .bind(&payment.subject) - .bind_timestamp(&payment.execution_time) - .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str())) - .bind(&payment.id.e2e_id) - .bind(&payment.id.msg_id) - .bind(&payment.id.sref) - .bind(subject.as_ref().map(|s| &s.wtid)) - .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str())) - .bind(subject.as_ref().map(|s| &s.metadata)) - .try_map(|r: PgRow| { - Ok(OutgoingRegistrationResult { - id: r.try_get_u64(0)?, - initiated: r.try_get_flag(1)?, - new: !r.try_get_flag(2)?, + ) + .bind(payment.amount) + .bind(payment.debit_fee) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str())) + .bind(&payment.id.e2e_id) + .bind(&payment.id.msg_id) + .bind(&payment.id.sref) + .bind(subject.as_ref().map(|s| &s.wtid)) + .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str())) + .bind(subject.as_ref().map(|s| &s.metadata)) + .try_map(|r: PgRow| { + Ok(OutgoingRegistrationResult { + id: r.try_get_u64(0)?, + initiated: r.try_get_flag(1)?, + new: !r.try_get_flag(2)?, + }) }) - }) - .fetch_one(pool) - .await + .fetch_one(pool) + ) } /// Register an outgoing batch @@ -71,32 +73,33 @@ pub async fn register_out_batch( payment: &OutTx, subject: Option<&OutgoingSubject>, ) -> sqlx::Result<OutgoingRegistrationResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_tx_id, out_initiated, out_found FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) ", - ) - .bind(payment.amount) - .bind(payment.debit_fee) - .bind(&payment.subject) - .bind_timestamp(&payment.execution_time) - .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str())) - .bind(&payment.id.e2e_id) - .bind(&payment.id.msg_id) - .bind(&payment.id.sref) - .bind(subject.as_ref().map(|s| &s.wtid)) - .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str())) - .bind(subject.as_ref().map(|s| &s.metadata)) - .try_map(|r: PgRow| { - Ok(OutgoingRegistrationResult { - id: r.try_get_u64(0)?, - initiated: r.try_get_flag(1)?, - new: !r.try_get_flag(2)?, + ) + .bind(payment.amount) + .bind(payment.debit_fee) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str())) + .bind(&payment.id.e2e_id) + .bind(&payment.id.msg_id) + .bind(&payment.id.sref) + .bind(subject.as_ref().map(|s| &s.wtid)) + .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str())) + .bind(subject.as_ref().map(|s| &s.metadata)) + .try_map(|r: PgRow| { + Ok(OutgoingRegistrationResult { + id: r.try_get_u64(0)?, + initiated: r.try_get_flag(1)?, + new: !r.try_get_flag(2)?, + }) }) - }) - .fetch_one(pool) - .await + .fetch_one(pool) + ) } #[derive(Debug, Clone, PartialEq, Eq)] @@ -119,31 +122,32 @@ pub enum IncomingRegistrationResult { /** Register an incoming payment */ pub async fn register_in(pool: &PgPool, payment: &InTx) -> sqlx::Result<InResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_found, out_completed, out_tx_id, out_bounce_id FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,NULL) ", - ) - .bind(payment.amount) - .bind(payment.credit_fee) - .bind(&payment.subject) - .bind_timestamp(&payment.execution_time) - .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) - .bind(payment.id.uetr) - .bind(&payment.id.tx_id) - .bind(&payment.id.sref) - .try_map(|r: PgRow| { - Ok(InResult { - id: r.try_get_u64("out_tx_id")?, - new: !r.try_get_flag("out_found")?, - completed: r.try_get_flag("out_completed")?, - bounce_id: r.try_get("out_bounce_id")?, - pending: false, + ) + .bind(payment.amount) + .bind(payment.credit_fee) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) + .bind(payment.id.uetr) + .bind(&payment.id.tx_id) + .bind(&payment.id.sref) + .try_map(|r: PgRow| { + Ok(InResult { + id: r.try_get_u64("out_tx_id")?, + new: !r.try_get_flag("out_found")?, + completed: r.try_get_flag("out_completed")?, + bounce_id: r.try_get("out_bounce_id")?, + pending: false, + }) }) - }) - .fetch_one(pool) - .await + .fetch_one(pool) + ) } /** Register an talerable incoming payment */ @@ -152,8 +156,9 @@ pub async fn register_in_talerable( payment: &InTx, subject: &IncomingSubject, ) -> sqlx::Result<IncomingRegistrationResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_reserve_pub_reuse, out_mapping_reuse, @@ -165,36 +170,36 @@ pub async fn register_in_talerable( out_bounce_id FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9::taler_incoming_type,$10,NULL) ", - ) - .bind(payment.amount) - .bind(payment.credit_fee) - .bind(&payment.subject) - .bind_timestamp(&payment.execution_time) - .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) - .bind(payment.id.uetr) - .bind(&payment.id.tx_id) - .bind(&payment.id.sref) - .bind(subject.ty()) - .bind(subject.key()) - .try_map(|r: PgRow| { - Ok(if r.try_get_flag("out_reserve_pub_reuse")? { - IncomingRegistrationResult::ReservePubReuse - } else if r.try_get_flag("out_mapping_reuse")? { - IncomingRegistrationResult::MappingReuse - } else if r.try_get_flag("out_unknown_mapping")? { - IncomingRegistrationResult::UnknownMapping - } else { - IncomingRegistrationResult::Success(InResult { - id: r.try_get_u64("out_tx_id")?, - new: !r.try_get_flag("out_found")?, - completed: r.try_get_flag("out_completed")?, - bounce_id: r.try_get("out_bounce_id")?, - pending: r.try_get("out_pending")?, + ) + .bind(payment.amount) + .bind(payment.credit_fee) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) + .bind(payment.id.uetr) + .bind(&payment.id.tx_id) + .bind(&payment.id.sref) + .bind(subject.ty().as_ref()) + .bind(subject.key()) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_reserve_pub_reuse")? { + IncomingRegistrationResult::ReservePubReuse + } else if r.try_get_flag("out_mapping_reuse")? { + IncomingRegistrationResult::MappingReuse + } else if r.try_get_flag("out_unknown_mapping")? { + IncomingRegistrationResult::UnknownMapping + } else { + IncomingRegistrationResult::Success(InResult { + id: r.try_get_u64("out_tx_id")?, + new: !r.try_get_flag("out_found")?, + completed: r.try_get_flag("out_completed")?, + bounce_id: r.try_get("out_bounce_id")?, + pending: r.try_get("out_pending")?, + }) }) }) - }) - .fetch_one(pool) - .await + .fetch_one(pool) + ) } /** Register an talerable incoming payment */ @@ -203,8 +208,9 @@ pub async fn register_in_qr_bill( payment: &InTx, reference: &str, ) -> sqlx::Result<IncomingRegistrationResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_reserve_pub_reuse, out_mapping_reuse, @@ -216,35 +222,35 @@ pub async fn register_in_qr_bill( out_bounce_id FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,$9) ", - ) - .bind(payment.amount) - .bind(payment.credit_fee) - .bind(&payment.subject) - .bind_timestamp(&payment.execution_time) - .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) - .bind(payment.id.uetr) - .bind(&payment.id.tx_id) - .bind(&payment.id.sref) - .bind(reference) - .try_map(|r: PgRow| { - Ok(if r.try_get_flag("out_reserve_pub_reuse")? { - IncomingRegistrationResult::ReservePubReuse - } else if r.try_get_flag("out_mapping_reuse")? { - IncomingRegistrationResult::MappingReuse - } else if r.try_get_flag("out_unknown_mapping")? { - IncomingRegistrationResult::UnknownMapping - } else { - IncomingRegistrationResult::Success(InResult { - id: r.try_get_u64("out_tx_id")?, - new: !r.try_get_flag("out_found")?, - completed: r.try_get_flag("out_completed")?, - bounce_id: r.try_get("out_bounce_id")?, - pending: r.try_get("out_pending")?, + ) + .bind(payment.amount) + .bind(payment.credit_fee) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) + .bind(payment.id.uetr) + .bind(&payment.id.tx_id) + .bind(&payment.id.sref) + .bind(reference) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_reserve_pub_reuse")? { + IncomingRegistrationResult::ReservePubReuse + } else if r.try_get_flag("out_mapping_reuse")? { + IncomingRegistrationResult::MappingReuse + } else if r.try_get_flag("out_unknown_mapping")? { + IncomingRegistrationResult::UnknownMapping + } else { + IncomingRegistrationResult::Success(InResult { + id: r.try_get_u64("out_tx_id")?, + new: !r.try_get_flag("out_found")?, + completed: r.try_get_flag("out_completed")?, + bounce_id: r.try_get("out_bounce_id")?, + pending: r.try_get("out_pending")?, + }) }) }) - }) - .fetch_one(pool) - .await + .fetch_one(pool) + ) } #[derive(Debug, Clone, PartialEq, Eq)] @@ -263,39 +269,40 @@ pub async fn register_in_malformed( timestamp: &Timestamp, cause: &str, ) -> sqlx::Result<IncomingBounceRegistrationResult> { - sqlx::query( - " + serialized!( + sqlx::query( + " SELECT out_found, out_tx_id, out_completed, out_bounce_id, out_talerable FROM register_and_bounce_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12) ", - ) - .bind(payment.amount) - .bind(payment.credit_fee) - .bind(&payment.subject) - .bind_timestamp(&payment.execution_time) - .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) - .bind(payment.id.uetr) - .bind(&payment.id.tx_id) - .bind(&payment.id.sref) - .bind(bounce_amount) - .bind_timestamp(timestamp) - .bind(bounce_end_to_end_id) - .bind(cause) - .try_map(|r: PgRow| { - Ok(if r.try_get_flag("out_talerable")? { - IncomingBounceRegistrationResult::Talerable - } else { - IncomingBounceRegistrationResult::Success(InResult { - id: r.try_get_u64("out_tx_id")?, - new: !r.try_get_flag("out_found")?, - completed: r.try_get_flag("out_completed")?, - bounce_id: r.try_get("out_bounce_id")?, - pending: false, + ) + .bind(payment.amount) + .bind(payment.credit_fee) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str())) + .bind(payment.id.uetr) + .bind(&payment.id.tx_id) + .bind(&payment.id.sref) + .bind(bounce_amount) + .bind_timestamp(timestamp) + .bind(bounce_end_to_end_id) + .bind(cause) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag("out_talerable")? { + IncomingBounceRegistrationResult::Talerable + } else { + IncomingBounceRegistrationResult::Success(InResult { + id: r.try_get_u64("out_tx_id")?, + new: !r.try_get_flag("out_found")?, + completed: r.try_get_flag("out_completed")?, + bounce_id: r.try_get("out_bounce_id")?, + pending: false, + }) }) }) - }) - .fetch_one(pool) - .await + .fetch_one(pool) + ) } #[cfg(test)] diff --git a/crates/libeufin-nexus/src/db/transfer.rs b/crates/libeufin-nexus/src/db/transfer.rs @@ -53,7 +53,7 @@ pub async fn transfer_register( ) ", ) - .bind(ty) + .bind(ty.as_ref()) .bind(account_pub) .bind(auth_pub) .bind(auth_sig) diff --git a/crates/libeufin-nexus/src/test.rs b/crates/libeufin-nexus/src/test.rs @@ -33,7 +33,6 @@ use taler_common::{ payto::{IbanPayto, PaytoURI, payto}, }, }; -use url::Url; use crate::{ config::{AccountType, NexusIngestCfg}, @@ -110,7 +109,7 @@ pub async fn gen_initiate( ) -> PaymentInitiationResult { let init = gen_init_pay(end_to_end_id, subject); initiate( - &db, + db, &init.amount, &init.subject, &init.creditor, @@ -142,7 +141,7 @@ async fn prepare(db: &PgPool) -> String { .await .unwrap() ); - return reference_number; + reference_number } /// Register a talerable reserve prepared incoming transaction @@ -282,7 +281,7 @@ pub async fn talerable_out(db: &PgPool) { db, &gen_out_pay(fmt_out_subject( &Base32::rand(), - &Url::from_str("https://exchange.test.com").unwrap(), + "https://exchange.test.com", None, )), )