commit 290a07f14a882208b25653c47e863ca9fba0504a
parent d6d7828fd9e9117796f62e772d8c3a917a394f7c
Author: Antoine A <>
Date: Fri, 24 Apr 2026 10:37:58 +0200
refactoring
Diffstat:
| M | src/db.rs | | | 2021 | +------------------------------------------------------------------------------ |
| A | src/db/initiated.rs | | | 886 | +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ |
| A | src/db/payment.rs | | | 1169 | +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ |
| M | src/lib.rs | | | 252 | ++----------------------------------------------------------------------------- |
| M | src/model.rs | | | 29 | +++++++++++++++++++++++++++++ |
| A | src/worker.rs | | | 245 | +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ |
6 files changed, 2358 insertions(+), 2244 deletions(-)
diff --git a/src/db.rs b/src/db.rs
@@ -14,16 +14,9 @@
TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/>
*/
-use std::collections::BTreeMap;
-
-use compact_str::CompactString;
-use const_format::formatcp;
use jiff::Timestamp;
use sqlx::{PgPool, QueryBuilder, Row, postgres::PgRow};
-use taler_api::{
- db::{BindHelper, IncomingType, TypeHelper, history},
- subject::{IncomingSubject, OutgoingSubject},
-};
+use taler_api::db::{BindHelper, IncomingType, TypeHelper, history};
use taler_common::{
api_common::{EddsaPublicKey, EddsaSignature},
api_params::History,
@@ -33,15 +26,14 @@ use taler_common::{
};
use tokio::sync::watch::{Receiver, Sender};
-use crate::{
- IncomingPayment, OutgoingPayment,
- config::parse_db_cfg,
- model::{InitiatedPayment, OutgoingId, PaymentBatch, SubmissionState},
-};
+use crate::config::parse_db_cfg;
+
+pub mod initiated;
+pub mod payment;
const SCHEMA: &str = "libeufin_nexus";
-const UNSETTLED: &str = "status NOT IN ('success', 'permanent_failure', 'late_failure')";
-const PENDING: &str = "status IN ('unsubmitted', 'pending')";
+pub const UNSETTLED: &str = "status NOT IN ('success', 'permanent_failure', 'late_failure')";
+pub const PENDING: &str = "status IN ('unsubmitted', 'pending')";
pub async fn pool(cfg: &Config) -> anyhow::Result<PgPool> {
let db = parse_db_cfg(cfg)?;
@@ -76,718 +68,6 @@ pub async fn notification_listener(
)
}
-/// Outgoing payments initiation result
-#[derive(Debug, PartialEq, Eq)]
-pub enum PaymentInitiationResult {
- Success(u64),
- RequestUidReuse,
-}
-
-/// Initiate a new payment
-pub async fn initiate(
- pool: &PgPool,
- payment: &InitiatedPayment,
-) -> sqlx::Result<PaymentInitiationResult> {
- let res = sqlx::query(
- "
- INSERT INTO initiated_outgoing_transactions (
- amount,
- subject,
- credit_payto,
- initiation_time,
- end_to_end_id
- ) VALUES ($1,$2,$3,$4,$5)
- RETURNING initiated_outgoing_transaction_id
- ",
- )
- .bind(&payment.amount)
- .bind(&payment.subject)
- .bind(payment.creditor.as_ref().as_str())
- .bind_timestamp(&payment.initiation_time)
- .bind(&payment.end_to_end_id)
- .try_map(|r: PgRow| Ok(PaymentInitiationResult::Success(r.try_get_u64(0)?)))
- .fetch_one(pool)
- .await;
- if let Err(e) = &res
- && let Some(db_err) = e.as_database_error()
- && db_err.code() == Some(std::borrow::Cow::Borrowed("23505"))
- {
- Ok(PaymentInitiationResult::RequestUidReuse)
- } else {
- res
- }
-}
-
-/// Group unbatched transaction into a single batch
-pub async fn batch_initiated(
- pool: &PgPool,
- timestamp: &Timestamp,
- ebics_id: &str,
- require_ack: bool,
-) -> sqlx::Result<()> {
- sqlx::query("SELECT batch_outgoing_transactions($1, $2, $3)")
- .bind_timestamp(timestamp)
- .bind(ebics_id)
- .bind(require_ack)
- .execute(pool)
- .await?;
- Ok(())
-}
-
-pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<()> {
- sqlx::query("UPDATE initiated_outgoing_transactions SET awaiting_ack=false WHERE initiated_outgoing_transaction_id=$1")
- .bind(id as i64)
- .execute(db)
- .await?;
- Ok(())
-}
-
-pub async fn initiated_submittable(
- db: &PgPool,
- currency: &Currency,
-) -> sqlx::Result<Vec<PaymentBatch>> {
- const SELECT_PART: &str = "
- 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!(
- "
- ({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(),
- })
- })
- .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
- ,subject
- ,credit_payto
- ,initiated_outgoing_transactions.initiation_time
- ,end_to_end_id
- ,initiated_outgoing_batch_id
- FROM initiated_outgoing_transactions
- JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
- WHERE initiated_outgoing_batches.status IN ('unsubmitted', 'transient_failure')
- ",
- )
- .try_map(|r: PgRow| {
- let payment = InitiatedPayment {
- 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")?,
- end_to_end_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)
-}
-
-pub async fn unsettled_tx_in_batch(
- db: &PgPool,
- currency: &Currency,
- msg_id: &str,
- execution_time: &Timestamp,
-) -> sqlx::Result<Vec<OutgoingPayment>> {
- sqlx::query(formatcp!(
- "
- SELECT
- end_to_end_id,
- amount,
- subject,
- credit_payto
- FROM initiated_outgoing_transactions
- JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
- WHERE message_id = $1
- AND initiated_outgoing_transactions.{UNSETTLED}
- "
- ))
- .bind(msg_id)
- .try_map(|r: PgRow| {
- Ok(OutgoingPayment {
- id: OutgoingId {
- msg_id: Some(msg_id.into()),
- end_to_end_id: r.try_get("end_to_end_id")?,
- acct_svcr_ref: None,
- },
- amount: r.try_get_amount("amount", currency)?,
- debit_fee: None,
- subject: r.try_get("subject")?,
- execution_time: *execution_time,
- creditor: r.try_get_opt_payto("credit_payto")?,
- })
- })
- .fetch_all(db)
- .await
-}
-
-/** Register submission success of order [orderId] for batch [id] at [timestamp] */
-pub async fn batch_sub_success(
- db: &PgPool,
- batch_id: u64,
- timestamp: &Timestamp,
- order_id: &str,
-) -> sqlx::Result<()> {
- 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
- ,order_id = $2
- ,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(batch_id as i64)
- .execute(&mut *tx)
- .await?;
- }
- tx.commit().await
-}
-
-/** Register submission failure with [msg] for batch [id] at [timestamp]*/
-pub async fn batch_sub_failure(
- db: &PgPool,
- batch_id: u64,
- timestamp: &Timestamp,
- msg: &str,
-) -> sqlx::Result<()> {
- let permanent = false;
- let mut tx = db.begin().await?;
- // Update batch status
- sqlx::query(
- "
- UPDATE initiated_outgoing_batches
- SET status = $1
- ,submission_date = $2
- ,status_msg = $3
- ,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
- 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(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
- SET status = 'pending', status_msg = $1
- WHERE order_id = $1 AND {PENDING}
- RETURNING initiated_outgoing_batch_id
- "
- ))
- .bind(msg)
- .bind(order_id)
- .try_map(|r: PgRow| r.try_get_u32(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)
- .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
- 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_u32(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))
-}
-
-/** Register order failure for [orderId] and return message_id and previous status_msg if found */
-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
- 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
- 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)))
-}
-
-/** Register payment status [state] with [msg] for batch [msgId] */
-pub async fn batch_status_update(
- db: &PgPool,
- msg_id: &str,
- 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
-}
-
-/** Register payment status [state] with [msg] for transaction [endToEndId] in batch [msgId] */
-pub async fn tx_status_update(
- db: &PgPool,
- end_to_end_id: &str,
- msg_id: &str,
- 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
-}
-
-#[derive(Debug, PartialEq, Eq)]
-pub struct OutgoingRegistrationResult {
- pub id: u64,
- pub initiated: bool,
- pub new: bool,
-}
-
-/** Register an outgoing payment reconciling it with its initiated payment counterpart if present */
-pub async fn register_out_tx(
- pool: &PgPool,
- payment: &OutgoingPayment,
- subject: Option<&OutgoingSubject>,
-) -> sqlx::Result<OutgoingRegistrationResult> {
- sqlx::query(
- "
- SELECT out_tx_id, out_initiated, out_found
- FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11)
- ",
- )
- .bind(&payment.amount)
- .bind(
- payment
- .debit_fee
- .as_ref()
- .unwrap_or(&Amount::zero(&payment.amount.currency)),
- )
- .bind(&payment.subject)
- .bind_timestamp(&payment.execution_time)
- .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str()))
- .bind(&payment.id.end_to_end_id)
- .bind(&payment.id.msg_id)
- .bind(&payment.id.acct_svcr_ref)
- .bind(subject.as_ref().map(|s| &s.wtid))
- .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str()))
- .bind(subject.as_ref().map(|s| &s.metadata))
- .try_map(|r: PgRow| {
- Ok(OutgoingRegistrationResult {
- id: r.try_get_u64(0)?,
- initiated: r.try_get_flag(1)?,
- new: !r.try_get_flag(2)?,
- })
- })
- .fetch_one(pool)
- .await
-}
-
-/// Register an outgoing batch
-pub async fn register_out_batch(
- pool: &PgPool,
- currency: &Currency,
- payment: &OutgoingPayment,
- subject: Option<&OutgoingSubject>,
-) -> sqlx::Result<OutgoingRegistrationResult> {
- sqlx::query(
- "
- SELECT out_tx_id, out_initiated, out_found
- FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11)
- ",
- )
- .bind(&payment.amount)
- .bind(
- payment
- .debit_fee
- .as_ref()
- .unwrap_or(&Amount::zero(currency)),
- )
- .bind(&payment.subject)
- .bind_timestamp(&payment.execution_time)
- .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str()))
- .bind(&payment.id.end_to_end_id)
- .bind(&payment.id.msg_id)
- .bind(&payment.id.acct_svcr_ref)
- .bind(subject.as_ref().map(|s| &s.wtid))
- .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str()))
- .bind(subject.as_ref().map(|s| &s.metadata))
- .try_map(|r: PgRow| {
- Ok(OutgoingRegistrationResult {
- id: r.try_get_u64(0)?,
- initiated: r.try_get_flag(1)?,
- new: !r.try_get_flag(2)?,
- })
- })
- .fetch_one(pool)
- .await
-}
-
-#[derive(Debug, Clone, PartialEq, Eq)]
-pub struct InResult {
- pub id: u64,
- pub new: bool,
- pub completed: bool,
- pub pending: bool,
- pub bounce_id: Option<CompactString>,
-}
-
-/** Incoming payments registration result */
-#[derive(Debug, PartialEq, Eq)]
-pub enum IncomingRegistrationResult {
- Success(InResult),
- ReservePubReuse,
- MappingReuse,
- UnknownMapping,
-}
-
-/** Register an incoming payment */
-pub async fn register_in(pool: &PgPool, payment: &IncomingPayment) -> sqlx::Result<InResult> {
- sqlx::query(
- "
- SELECT out_found, out_completed, out_tx_id, out_bounce_id
- FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,NULL)
- ",
- )
- .bind(&payment.amount)
- .bind(
- payment
- .credit_fee
- .as_ref()
- .unwrap_or(&Amount::zero(&payment.amount.currency)),
- )
- .bind(&payment.subject)
- .bind_timestamp(&payment.execution_time)
- .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
- .bind(payment.id.uetr)
- .bind(&payment.id.tx_id)
- .bind(&payment.id.acct_svcr_ref)
- .try_map(|r: PgRow| {
- Ok(InResult {
- id: r.try_get_u64("out_tx_id")?,
- new: !r.try_get_flag("out_found")?,
- completed: r.try_get_flag("out_completed")?,
- bounce_id: r.try_get("out_bounce_id")?,
- pending: false,
- })
- })
- .fetch_one(pool)
- .await
-}
-
-/** Register an talerable incoming payment */
-pub async fn register_in_talerable(
- pool: &PgPool,
- payment: &IncomingPayment,
- subject: &IncomingSubject,
-) -> sqlx::Result<IncomingRegistrationResult> {
- sqlx::query(
- "
- SELECT
- out_reserve_pub_reuse,
- out_mapping_reuse,
- out_unknown_mapping,
- out_found,
- out_completed,
- out_pending,
- out_tx_id,
- out_bounce_id
- FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9::taler_incoming_type,$10,NULL)
- ",
- )
- .bind(&payment.amount)
- .bind(
- payment
- .credit_fee
- .as_ref()
- .unwrap_or(&Amount::zero(&payment.amount.currency)),
- )
- .bind(&payment.subject)
- .bind_timestamp(&payment.execution_time)
- .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
- .bind(payment.id.uetr)
- .bind(&payment.id.tx_id)
- .bind(&payment.id.acct_svcr_ref)
- .bind(subject.ty())
- .bind(subject.key())
- .try_map(|r: PgRow| {
- Ok(if r.try_get_flag("out_reserve_pub_reuse")? {
- IncomingRegistrationResult::ReservePubReuse
- } else if r.try_get_flag("out_mapping_reuse")? {
- IncomingRegistrationResult::MappingReuse
- } else if r.try_get_flag("out_unknown_mapping")? {
- IncomingRegistrationResult::UnknownMapping
- } else {
- IncomingRegistrationResult::Success(InResult {
- id: r.try_get_u64("out_tx_id")?,
- new: !r.try_get_flag("out_found")?,
- completed: r.try_get_flag("out_completed")?,
- bounce_id: r.try_get("out_bounce_id")?,
- pending: r.try_get("out_pending")?,
- })
- })
- })
- .fetch_one(pool)
- .await
-}
-
-/** Register an talerable incoming payment */
-pub async fn register_in_qr_bill(
- pool: &PgPool,
- payment: &IncomingPayment,
- reference: &str,
-) -> sqlx::Result<IncomingRegistrationResult> {
- sqlx::query(
- "
- SELECT
- out_reserve_pub_reuse,
- out_mapping_reuse,
- out_unknown_mapping,
- out_found,
- out_completed,
- out_pending,
- out_tx_id,
- out_bounce_id
- FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,$9)
- ",
- )
- .bind(&payment.amount)
- .bind(
- payment
- .credit_fee
- .as_ref()
- .unwrap_or(&Amount::zero(&payment.amount.currency)),
- )
- .bind(&payment.subject)
- .bind_timestamp(&payment.execution_time)
- .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
- .bind(payment.id.uetr)
- .bind(&payment.id.tx_id)
- .bind(&payment.id.acct_svcr_ref)
- .bind(reference)
- .try_map(|r: PgRow| {
- Ok(if r.try_get_flag("out_reserve_pub_reuse")? {
- IncomingRegistrationResult::ReservePubReuse
- } else if r.try_get_flag("out_mapping_reuse")? {
- IncomingRegistrationResult::MappingReuse
- } else if r.try_get_flag("out_unknown_mapping")? {
- IncomingRegistrationResult::UnknownMapping
- } else {
- IncomingRegistrationResult::Success(InResult {
- id: r.try_get_u64("out_tx_id")?,
- new: !r.try_get_flag("out_found")?,
- completed: r.try_get_flag("out_completed")?,
- bounce_id: r.try_get("out_bounce_id")?,
- pending: r.try_get("out_pending")?,
- })
- })
- })
- .fetch_one(pool)
- .await
-}
-
-#[derive(Debug, Clone, PartialEq, Eq)]
-/** Incoming payments bounce registration result */
-pub enum IncomingBounceRegistrationResult {
- Success(InResult),
- Talerable,
-}
-
-/** Register an incoming payment and bounce it */
-pub async fn register_in_malformed(
- pool: &PgPool,
- payment: &IncomingPayment,
- bounce_amount: &Amount,
- bounce_end_to_end_id: &str,
- timestamp: &Timestamp,
- cause: &str,
-) -> sqlx::Result<IncomingBounceRegistrationResult> {
- sqlx::query(
- "
- SELECT out_found, out_tx_id, out_completed, out_bounce_id, out_talerable
- FROM register_and_bounce_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12)
- ",
- )
- .bind(&payment.amount)
- .bind(
- payment
- .credit_fee
- .as_ref()
- .unwrap_or(&Amount::zero(&payment.amount.currency)),
- )
- .bind(&payment.subject)
- .bind_timestamp(&payment.execution_time)
- .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
- .bind(payment.id.uetr)
- .bind(&payment.id.tx_id)
- .bind(&payment.id.acct_svcr_ref)
- .bind(bounce_amount)
- .bind_timestamp(timestamp)
- .bind(bounce_end_to_end_id)
- .bind(cause)
- .try_map(|r: PgRow| {
- Ok(if r.try_get_flag("out_talerable")? {
- IncomingBounceRegistrationResult::Talerable
- } else {
- IncomingBounceRegistrationResult::Success(InResult {
- id: r.try_get_u64("out_tx_id")?,
- new: !r.try_get_flag("out_found")?,
- completed: r.try_get_flag("out_completed")?,
- bounce_id: r.try_get("out_bounce_id")?,
- pending: false,
- })
- })
- })
- .fetch_one(pool)
- .await
-}
-
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RegistrationResult {
Success,
@@ -972,51 +252,29 @@ pub async fn incoming_history(
#[cfg(test)]
mod test {
- use std::str::FromStr;
- use std::sync::LazyLock;
+ use std::{str::FromStr, sync::LazyLock};
use compact_str::CompactString;
- use jiff::{Span, Timestamp, civil::Date};
- use sqlx::PgPool;
- use sqlx::{PgConnection, Row, postgres::PgRow};
- use taler_api::{
- db::{IncomingType, TypeHelper},
- subject::subject_fmt_qr_bill,
- };
- use taler_common::{
- api_common::ShortHashCode,
- types::{amount::Amount, utils::date_to_utc_ts},
- };
+ use jiff::Timestamp;
+ use sqlx::{PgConnection, PgPool, Row, postgres::PgRow};
+ use taler_api::db::{IncomingType, TypeHelper};
use taler_common::{
- api_common::{EddsaPublicKey, EddsaSignature},
- types::amount::amount,
- };
- use taler_common::{
- config::Config,
- types::{amount::Currency, payto::IbanPayto},
- };
- use uuid::Uuid;
-
- use crate::{
- CONFIG_SOURCE, IncomingPayment, OutgoingBatch, Tx,
- config::{AccountType, NexusCfg, NexusIngestConfig},
- db::{
- InResult, IncomingBounceRegistrationResult, OutgoingRegistrationResult,
- RegistrationResult, batch_status_update, batch_sub_failure, batch_sub_success,
- initiated_submittable, order_failure, order_step, order_success, register_in_malformed,
- transfer_register, tx_status_update,
+ api_common::EddsaPublicKey,
+ types::{
+ amount::{Amount, Currency},
+ payto::IbanPayto,
},
- model::{IncomingId, InitiatedPayment, OutgoingId, SubmissionState},
- register_incoming, register_tx,
};
- use crate::{OutgoingPayment, rand_ebics_id, register_outgoing, register_outgoing_batch};
+
use crate::{
- db::PaymentInitiationResult, db::batch_initiated, db::initiate, db::initiated_ack,
+ CONFIG_SOURCE,
+ model::{IncomingId, IncomingPayment, InitiatedPayment, OutgoingId, OutgoingPayment},
+ rand_ebics_id,
};
pub static CURRENCY: LazyLock<Currency> = LazyLock::new(|| "KUDOS".parse().unwrap());
- async fn setup() -> (PgConnection, PgPool) {
+ 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)
@@ -1075,7 +333,7 @@ mod test {
}
}
- async fn check_count(db: &PgPool, nb_tx: usize, nb_bounce: usize) {
+ pub async fn check_count(db: &PgPool, nb_tx: usize, nb_bounce: usize) {
sqlx::query(
"
SELECT (SELECT count(*) FROM incoming_transactions) + (SELECT count(*) FROM outgoing_transactions) AS transactions,
@@ -1094,7 +352,7 @@ mod test {
.unwrap();
}
- async fn check_in_count(
+ pub async fn check_in_count(
db: &PgPool,
nb_incoming: usize,
nb_bounce: usize,
@@ -1119,7 +377,7 @@ mod test {
.unwrap();
}
- async fn check_out_count(db: &PgPool, nb_outgoing: u64, nb_talerable: u64) {
+ pub async fn check_out_count(db: &PgPool, nb_outgoing: u64, nb_talerable: u64) {
sqlx::query(
"
SELECT (SELECT count(*) FROM outgoing_transactions) AS outgoing,
@@ -1139,7 +397,7 @@ mod test {
}
#[derive(Debug, PartialEq, Eq)]
- enum Status {
+ pub enum Status {
Simple,
Pending,
Bounced,
@@ -1148,9 +406,7 @@ mod test {
Kyc(EddsaPublicKey),
}
- use Status::*;
-
- async fn check_in(db: &PgPool, state: &[Status]) {
+ pub async fn check_in(db: &PgPool, state: &[Status]) {
let current = sqlx::query(
"
SELECT pending_recurrent_incoming_transactions.authorization_pub IS NOT NULL, initiated_outgoing_transaction_id IS NOT NULL, debit_payto IS NULL OR subject IS NULL, type::text, metadata
@@ -1184,1231 +440,4 @@ mod test {
.unwrap();
assert_eq!(state, current);
}
-
- #[tokio::test]
- async fn out_tx() {
- let (_, db) = setup().await;
- // Register initiated transactions
- for subject in [
- "initiated by nexus".to_owned(),
- format!("{} https://exchange.com/", ShortHashCode::rand()),
- ] {
- let payment = gen_out_pay(subject.clone());
- assert!(matches!(
- initiate(
- &db,
- &gen_init_pay(payment.id.end_to_end_id.clone().unwrap(), subject),
- )
- .await,
- Ok(PaymentInitiationResult::Success(_))
- ));
- let first = register_outgoing(&db, &payment).await.unwrap();
- assert_eq!(
- first,
- OutgoingRegistrationResult {
- id: first.id,
- initiated: true,
- new: true
- }
- );
- assert_eq!(
- register_outgoing(&db, &payment).await.unwrap(),
- OutgoingRegistrationResult {
- id: first.id,
- initiated: true,
- new: false
- }
- );
- let payment = OutgoingPayment {
- id: OutgoingId {
- msg_id: None,
- end_to_end_id: None,
- acct_svcr_ref: payment.id.end_to_end_id,
- },
- ..payment
- };
- let second = register_outgoing(&db, &payment).await.unwrap();
- assert_eq!(
- second,
- OutgoingRegistrationResult {
- id: first.id + 1,
- initiated: false,
- new: true
- }
- );
- assert_eq!(
- register_outgoing(&db, &payment).await.unwrap(),
- OutgoingRegistrationResult {
- id: second.id,
- initiated: false,
- new: false
- }
- );
- }
- check_out_count(&db, 4, 1).await;
-
- // Register unknown
- for subject in [
- "initiated by nexus".to_owned(),
- format!("{} https://exchange.com/", ShortHashCode::rand()),
- ] {
- let payment = gen_out_pay(subject.clone());
- let res = register_outgoing(&db, &payment).await.unwrap();
- assert_eq!(
- res,
- OutgoingRegistrationResult {
- id: res.id,
- initiated: false,
- new: true
- }
- );
- assert_eq!(
- register_outgoing(&db, &payment).await.unwrap(),
- OutgoingRegistrationResult {
- id: res.id,
- initiated: false,
- new: false
- }
- );
- }
- check_out_count(&db, 6, 2).await;
-
- // Register wtid reuse
- let wtid = ShortHashCode::rand();
- for subject in [
- format!("{wtid} https://exchange.com/"),
- format!("{wtid} https://exchange.com/"),
- ] {
- let payment = gen_out_pay(subject.clone());
- let res = register_outgoing(&db, &payment).await.unwrap();
- assert_eq!(
- res,
- OutgoingRegistrationResult {
- id: res.id,
- initiated: false,
- new: true
- }
- );
- assert_eq!(
- register_outgoing(&db, &payment).await.unwrap(),
- OutgoingRegistrationResult {
- id: res.id,
- initiated: false,
- new: false
- }
- );
- }
- check_out_count(&db, 8, 3).await
- }
-
- #[tokio::test]
- async fn out_batch() {
- let (_, db) = setup().await;
- // Init batch
- let wtid = ShortHashCode::rand();
- for subject in [
- "initiated by nexus".to_string(),
- format!("{} https://exchange.com/", ShortHashCode::rand()),
- format!("{wtid} https://exchange.com/"),
- format!("{wtid} https://exchange.com/"),
- ] {
- assert!(matches!(
- initiate(&db, &gen_init_pay(rand_ebics_id(), subject),).await,
- Ok(PaymentInitiationResult::Success(_))
- ));
- }
- batch_initiated(&db, &Timestamp::now(), "BATCH", false)
- .await
- .unwrap();
-
- // Register batch
- register_outgoing_batch(
- &db,
- &CURRENCY,
- &OutgoingBatch {
- msg_id: "BATCH".into(),
- execution_time: Timestamp::now(),
- },
- )
- .await
- .unwrap();
- check_out_count(&db, 4, 2).await;
-
- // Test manual ack
- let mut txs = Vec::new();
- for nb in 0..3 {
- let res = initiate(&db, &gen_init_pay(rand_ebics_id(), format!("tx {nb}"))).await;
- if let Ok(PaymentInitiationResult::Success(id)) = &res {
- txs.push(*id);
- } else {
- panic!("Expected success got {res:?}");
- }
- }
-
- // Check not sent without ack
- batch_initiated(&db, &Timestamp::now(), "BATCH_MANUAL", true)
- .await
- .unwrap();
- register_outgoing_batch(
- &db,
- &CURRENCY,
- &OutgoingBatch {
- msg_id: "BATCH_MANUAL".into(),
- execution_time: Timestamp::now(),
- },
- )
- .await
- .unwrap();
- check_out_count(&db, 4, 2).await;
-
- // Check sent with ack
- for tx in txs {
- initiated_ack(&db, tx).await.unwrap();
- }
- batch_initiated(&db, &Timestamp::now(), "BATCH_MANUAL", true)
- .await
- .unwrap();
- register_outgoing_batch(
- &db,
- &CURRENCY,
- &OutgoingBatch {
- msg_id: "BATCH_MANUAL".into(),
- execution_time: Timestamp::now(),
- },
- )
- .await
- .unwrap();
- check_out_count(&db, 7, 2).await;
- }
-
- #[tokio::test]
- async fn in_bounce() {
- let (_, db) = setup().await;
-
- // Creating and bouncing one incoming transaction
- let payment = gen_in_pay("incoming and bounce");
- let id = rand_ebics_id();
-
- let bounce_amount = amount("KUDOS:2.53");
- let res = register_in_malformed(
- &db,
- &payment,
- &bounce_amount,
- &id,
- &Timestamp::now(),
- "manual bounce",
- )
- .await
- .unwrap();
- assert!(
- matches!(
- res,
- IncomingBounceRegistrationResult::Success(InResult {
- new: true,
- id: _,
- completed: false,
- pending: false,
- ref bounce_id
- }) if bounce_id.as_ref() == Some(&id)
- ),
- "{res:?}"
- );
- // Idempotent
- let res = register_in_malformed(
- &db,
- &payment,
- &amount("KUDOS:2.5"),
- &rand_ebics_id(),
- &Timestamp::now(),
- "other reason to bounce",
- )
- .await
- .unwrap();
- assert!(
- matches!(
- res,
- IncomingBounceRegistrationResult::Success(InResult {
- new: false,
- id: _,
- completed: false,
- pending: false,
- ref bounce_id
- }) if bounce_id.as_ref() == Some(&id)
- ),
- "{res:?}"
- );
-
- // Checking one incoming got created and bounced
- sqlx::query(
- "
- SELECT
- incoming_transactions.amount as in_amount,
- initiated_outgoing_transactions.amount as bounce_amount
- FROM incoming_transactions
- JOIN bounced_transactions USING (incoming_transaction_id)
- JOIN initiated_outgoing_transactions USING (initiated_outgoing_transaction_id)
- ",
- )
- .try_map(|r: PgRow| {
- assert_eq!(r.try_get_amount("in_amount", &CURRENCY)?, payment.amount);
- assert_eq!(r.try_get_amount("bounce_amount", &CURRENCY)?, bounce_amount);
- Ok(())
- })
- .fetch_one(&db)
- .await
- .unwrap();
- }
-
- #[tokio::test]
- async fn in_simple() {
- let (_, db) = setup().await;
-
- let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
-
- // Register
- let incoming = gen_in_pay("test".to_owned());
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Bounced]).await;
-
- // Idempotent
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Bounced]).await;
-
- // Many
- register_incoming(&db, &cfg, &gen_in_pay("another subject".to_owned()))
- .await
- .unwrap();
- check_in(&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;
-
- let original = gen_in_pay("test 2".to_owned());
- let incomplete = IncomingPayment {
- subject: None,
- debtor: None,
- ..original.clone()
- };
-
- // Register incomplete transaction
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await;
- // Idempotent
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await;
- // Recover info when completed
- register_incoming(&db, &cfg, &original).await.unwrap();
- check_in(&db, &[Bounced, Bounced, Simple, Bounced]).await;
- }
-
- #[tokio::test]
- async fn in_talerable() {
- let (_, db) = setup().await;
-
- let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
- let key = EddsaPublicKey::rand();
- let subject = format!("test with {key} reserve pub");
-
- // Register
- let incoming = gen_in_pay(subject.clone());
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Reserve(key.clone())]).await;
-
- // Idempotent
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Reserve(key.clone())]).await;
-
- // Key reuse is bounced
- register_incoming(&db, &cfg, &gen_in_pay(subject.clone()))
- .await
- .unwrap();
- register_incoming(&db, &cfg, &gen_in_pay(format!("another {subject}")))
- .await
- .unwrap();
- check_in(&db, &[Reserve(key.clone()), Bounced, Bounced]).await;
-
- // Admin balance adjust is ignored
- register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
- .await
- .unwrap();
- check_in(&db, &[Reserve(key.clone()), Bounced, Bounced, Simple]).await;
-
- let new = EddsaPublicKey::rand();
- let original = gen_in_pay(format!("test 2 with {new} reserve pub"));
- let incomplete = IncomingPayment {
- subject: None,
- debtor: None,
- ..original.clone()
- };
-
- // Register incomplete transaction
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(
- &db,
- &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete],
- )
- .await;
- // Idempotent
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(
- &db,
- &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete],
- )
- .await;
- // Recover info when completed
- register_incoming(&db, &cfg, &original).await.unwrap();
- check_in(
- &db,
- &[Reserve(key.clone()), Bounced, Bounced, Simple, Reserve(new)],
- )
- .await;
- }
-
- #[tokio::test]
- async fn in_mapping() {
- let (_, db) = setup().await;
- let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
- let first = EddsaPublicKey::rand();
- let auth_pub = EddsaPublicKey::rand();
- let auth_sig = EddsaSignature::rand();
- let reference_number = subject_fmt_qr_bill(auth_pub.as_slice());
- let subject = format!("test with MAP:{auth_pub} auth pub");
-
- assert_eq!(
- transfer_register(
- &db,
- IncomingType::reserve,
- &first,
- &auth_pub,
- &auth_sig,
- false,
- &reference_number,
- &Timestamp::now()
- )
- .await
- .unwrap(),
- RegistrationResult::Success
- );
-
- // Register
- let incoming = gen_in_pay(subject.clone());
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Reserve(first.clone())]).await;
-
- // Idempotent
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Reserve(first.clone())]).await;
-
- // Admin balance adjust is ignored
- register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
- .await
- .unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple]).await;
-
- let original = gen_in_pay(format!("test 2 for {subject}"));
- let incomplete = IncomingPayment {
- subject: None,
- debtor: None,
- ..original.clone()
- };
- // Register incomplete transaction
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
- // Idempotent
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
- // Recover info when completed
- register_incoming(&db, &cfg, &original).await.unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
-
- let second = EddsaPublicKey::rand();
- assert_eq!(
- transfer_register(
- &db,
- IncomingType::reserve,
- &second,
- &auth_pub,
- &auth_sig,
- true,
- &reference_number,
- &Timestamp::now()
- )
- .await
- .unwrap(),
- RegistrationResult::Success
- );
- check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
-
- // Key reuse is pending
- for _ in 0..3 {
- register_incoming(&db, &cfg, &gen_in_pay(subject.clone()))
- .await
- .unwrap();
- }
- check_in(
- &db,
- &[
- Reserve(first.clone()),
- Simple,
- Bounced,
- Reserve(second.clone()),
- Pending,
- Pending,
- ],
- )
- .await;
-
- // Finish pending
- let third = EddsaPublicKey::rand();
- assert_eq!(
- transfer_register(
- &db,
- IncomingType::reserve,
- &third,
- &auth_pub,
- &auth_sig,
- true,
- &reference_number,
- &Timestamp::now()
- )
- .await
- .unwrap(),
- RegistrationResult::Success
- );
- check_in(
- &db,
- &[
- Reserve(first.clone()),
- Simple,
- Bounced,
- Reserve(second.clone()),
- Reserve(third.clone()),
- Pending,
- ],
- )
- .await;
- }
-
- #[tokio::test]
- async fn in_reference() {
- let (_, db) = setup().await;
- let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
- let first = EddsaPublicKey::rand();
- let auth_pub = EddsaPublicKey::rand();
- let auth_sig = EddsaSignature::rand();
- let reference_number = subject_fmt_qr_bill(auth_pub.as_slice());
-
- assert_eq!(
- transfer_register(
- &db,
- IncomingType::reserve,
- &first,
- &auth_pub,
- &auth_sig,
- false,
- &reference_number,
- &Timestamp::now()
- )
- .await
- .unwrap(),
- RegistrationResult::Success
- );
-
- // Register
- let incoming = gen_in_pay(reference_number.clone());
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Reserve(first.clone())]).await;
-
- // Idempotent
- register_incoming(&db, &cfg, &incoming).await.unwrap();
- check_in(&db, &[Reserve(first.clone())]).await;
-
- // Admin balance adjust is ignored
- register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
- .await
- .unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple]).await;
-
- let original = gen_in_pay(reference_number.clone());
- let incomplete = IncomingPayment {
- subject: None,
- debtor: None,
- ..original.clone()
- };
- // Register incomplete transaction
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
- // Idempotent
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
- // Recover info when completed
- register_incoming(&db, &cfg, &original).await.unwrap();
- check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
-
- let second = EddsaPublicKey::rand();
- assert_eq!(
- transfer_register(
- &db,
- IncomingType::reserve,
- &second,
- &auth_pub,
- &auth_sig,
- true,
- &reference_number,
- &Timestamp::now()
- )
- .await
- .unwrap(),
- RegistrationResult::Success
- );
- check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
-
- // Key reuse is pending
- for _ in 0..3 {
- register_incoming(&db, &cfg, &gen_in_pay(reference_number.clone()))
- .await
- .unwrap();
- }
- check_in(
- &db,
- &[
- Reserve(first.clone()),
- Simple,
- Bounced,
- Reserve(second.clone()),
- Pending,
- Pending,
- ],
- )
- .await;
-
- // Finish pending
- let third = EddsaPublicKey::rand();
- assert_eq!(
- transfer_register(
- &db,
- IncomingType::reserve,
- &third,
- &auth_pub,
- &auth_sig,
- true,
- &reference_number,
- &Timestamp::now()
- )
- .await
- .unwrap(),
- RegistrationResult::Success
- );
- check_in(
- &db,
- &[
- Reserve(first.clone()),
- Simple,
- Bounced,
- Reserve(second.clone()),
- Reserve(third.clone()),
- Pending,
- ],
- )
- .await;
- }
-
- #[tokio::test]
- async fn in_recover_info() {
- let (_, db) = setup().await;
- let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
-
- async fn check_content(db: &PgPool, p: &IncomingPayment) {
- sqlx::query(
- "
- SELECT
- uetr IS NOT DISTINCT FROM $1 AND
- tx_id IS NOT DISTINCT FROM $2 AND
- acct_svcr_ref IS NOT DISTINCT FROM $3 AND
- subject IS NOT DISTINCT FROM $4 AND
- debit_payto IS NOT DISTINCT FROM $5
- FROM incoming_transactions ORDER BY incoming_transaction_id DESC LIMIT 1
- ",
- )
- .bind(p.id.uetr)
- .bind(&p.id.tx_id)
- .bind(&p.id.acct_svcr_ref)
- .bind(&p.subject)
- .bind(p.debtor.as_ref().map(|it| it.as_ref().as_str()))
- .try_map(|r: PgRow| {
- assert!(r.try_get_flag(0)?);
- Ok(())
- })
- .fetch_one(db)
- .await
- .unwrap();
- }
-
- // Non talerable
- for (i, id) in [
- IncomingId::new(Some(Uuid::new_v4()), None, None),
- IncomingId::new(None, Some(rand_ebics_id()), None),
- IncomingId::new(None, None, Some(rand_ebics_id())),
- ]
- .iter()
- .enumerate()
- {
- let payment = gen_in_pay("subject".to_owned());
-
- // Register minimal
- let partial = IncomingPayment {
- id: id.clone(),
- subject: None,
- debtor: None,
- ..payment.clone()
- };
- register_incoming(&db, &cfg, &partial).await.unwrap();
- check_content(&db, &partial).await;
- check_in_count(&db, i + 1, i, 0).await;
-
- // Recover ID
- let full_id = IncomingId::new(
- Some(id.uetr.unwrap_or_else(Uuid::new_v4)),
- Some(id.tx_id.clone().unwrap_or_else(rand_ebics_id)),
- Some(id.acct_svcr_ref.clone().unwrap_or_else(rand_ebics_id)),
- );
- let full = IncomingPayment {
- id: full_id.clone(),
- ..partial.clone()
- };
- register_incoming(&db, &cfg, &full).await.unwrap();
- check_content(&db, &full).await;
- check_in_count(&db, i + 1, i, 0).await;
-
- // Recover subject & debtor
- let full = IncomingPayment {
- id: full_id,
- ..payment.clone()
- };
- register_incoming(&db, &cfg, &full).await.unwrap();
- check_content(&db, &full).await;
- check_in_count(&db, i + 1, i + 1, 0).await;
- }
-
- // Talerable
- for (i, id) in [
- IncomingId::new(Some(Uuid::new_v4()), None, None),
- IncomingId::new(None, Some(rand_ebics_id()), None),
- IncomingId::new(None, None, Some(rand_ebics_id())),
- ]
- .iter()
- .enumerate()
- {
- let key = EddsaPublicKey::rand();
- let payment = gen_in_pay(format!("test with {key} reserve pub"));
-
- // Register minimal
- let partial = IncomingPayment {
- id: id.clone(),
- subject: None,
- debtor: None,
- ..payment.clone()
- };
- register_incoming(&db, &cfg, &partial).await.unwrap();
- check_content(&db, &partial).await;
- check_in_count(&db, i + 4, 3, i).await;
-
- // Recover ID
- let full_id = IncomingId::new(
- Some(id.uetr.unwrap_or_else(Uuid::new_v4)),
- Some(id.tx_id.clone().unwrap_or_else(rand_ebics_id)),
- Some(id.acct_svcr_ref.clone().unwrap_or_else(rand_ebics_id)),
- );
- let full = IncomingPayment {
- id: full_id.clone(),
- ..partial.clone()
- };
- register_incoming(&db, &cfg, &full).await.unwrap();
- check_content(&db, &full).await;
- check_in_count(&db, i + 4, 3, i).await;
-
- // Recover subject & debtor
- let full = IncomingPayment {
- id: full_id,
- ..payment.clone()
- };
- register_incoming(&db, &cfg, &full).await.unwrap();
- check_content(&db, &full).await;
- check_in_count(&db, i + 4, 3, i + 1).await;
- }
- }
-
- #[tokio::test]
- pub async fn in_horror() {
- let (_, db) = setup().await;
- let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
-
- // Check we do not bounce already registered talerable transaction
- let key = EddsaPublicKey::rand();
- let payment = gen_in_pay(format!("test with {key} reserve pub"));
- register_incoming(&db, &cfg, &payment).await.unwrap();
- assert_eq!(
- register_in_malformed(
- &db,
- &payment,
- &amount("KUDOS:2.53"),
- &rand_ebics_id(),
- &Timestamp::now(),
- "manual bounce",
- )
- .await
- .unwrap(),
- IncomingBounceRegistrationResult::Talerable
- );
- let incomplete = IncomingPayment {
- subject: None,
- ..payment.clone()
- };
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- register_incoming(&db, &cfg, &payment).await.unwrap();
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- check_in(&db, &[Reserve(key.clone())]).await;
-
- // Check we do not register as talerable bounced transaction
- let new_key = EddsaPublicKey::rand();
- let payment = gen_in_pay(format!("bounced {new_key}"));
- let incomplete = IncomingPayment {
- subject: None,
- ..payment.clone()
- };
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- register_incoming(&db, &cfg, &payment).await.unwrap();
- register_incoming(&db, &cfg, &incomplete).await.unwrap();
- register_incoming(&db, &cfg, &payment).await.unwrap();
- check_in(&db, &[Reserve(key.clone()), Bounced]).await;
- }
-
- #[tokio::test]
- pub async fn initiated_skip() {
- let (_, db) = setup().await;
- let cfg = Config::from_file(CONFIG_SOURCE, Some("libeufin-nexus/conf/skip.conf")).unwrap();
- let cfg = NexusCfg::parse(cfg).unwrap();
- let cfg = cfg.ingest().unwrap();
- let millis = Span::new().milliseconds(10);
-
- async fn ingest(db: &PgPool, cfg: &NexusIngestConfig, execution_time: Timestamp) {
- for tx in [
- Tx::In(
- gen_in_pay(format!("test at {execution_time}"))
- .with_execution_time(execution_time),
- ),
- Tx::Out(
- gen_out_pay(format!("test at {execution_time}"))
- .with_execution_time(execution_time),
- ),
- ] {
- register_tx(db, cfg, &tx).await.unwrap()
- }
- }
-
- assert_eq!(
- cfg.ignore_txs_before,
- date_to_utc_ts(&Date::from_str("2024-04-04").unwrap())
- );
- assert_eq!(
- cfg.ignore_bounces_before,
- date_to_utc_ts(&Date::from_str("2024-06-12").unwrap())
- );
-
- // No transaction at the beginning
- check_count(&db, 0, 0).await;
-
- // Skipped transactions
- ingest(&db, &cfg, cfg.ignore_txs_before - millis).await;
- check_count(&db, 0, 0).await;
-
- // Skipped bounces
- ingest(&db, &cfg, cfg.ignore_txs_before).await;
- ingest(&db, &cfg, cfg.ignore_txs_before + millis).await;
- ingest(&db, &cfg, cfg.ignore_bounces_before - millis).await;
- check_count(&db, 6, 0).await;
-
- // Bounces
- ingest(&db, &cfg, cfg.ignore_bounces_before).await;
- ingest(&db, &cfg, cfg.ignore_bounces_before + millis).await;
- check_count(&db, 10, 2).await;
- }
-
- #[tokio::test]
- pub async fn initiated_status() {
- use SubmissionState::*;
-
- let (_, db) = setup().await;
-
- let check_parts = async |batch_id: u64,
- batch_status: SubmissionState,
- batch_msg: &str,
- tx_status: SubmissionState,
- tx_msg: &str,
- settled_status: SubmissionState,
- settled_msg: &str| {
- // Check batch status
- let msg_id: String = sqlx::query(
- "
- SELECT message_id, status, status_msg FROM initiated_outgoing_batches WHERE initiated_outgoing_batch_id=?
- "
- ).bind(batch_id as i64)
- .try_map(|r: PgRow| {
- let msg_id: String = r.try_get("message_id")?;
- assert_eq!((batch_status, batch_msg), (r.try_get("status")?, r.try_get("status_msg")?), "{msg_id}");
- Ok(msg_id)
- }).fetch_one(&db).await.unwrap();
- // Check tx status
- sqlx::query(
- "
- SELECT end_to_end_id, status, status_msg FROM initiated_outgoing_transactions WHERE initiated_outgoing_batch_id=?
- "
- ).bind(batch_id as i64).try_map(|r: PgRow| {
- let end_to_end_id: &str = r.try_get("end_to_end_id")?;
- let expected = match end_to_end_id {
- "TX" => (tx_status, tx_msg),
- "TX_SETTLED" => (settled_status, settled_msg),
- _ =>panic!("Unexpected tx $endToEndId")
- };
- assert_eq!(expected,
- (r.try_get("status")?, r.try_get("status_msg")?),
- "{msg_id},{end_to_end_id}"
- );
- Ok(())
- }).fetch_all(&db).await.unwrap();
- };
-
- let check_batch_tx = async |batch_id: u64,
- status: SubmissionState,
- msg: &str,
- tx_status: SubmissionState| {
- check_parts(batch_id, status, msg, tx_status, msg, tx_status, msg).await;
- };
- let check_batch = async |batch_id: u64, status: SubmissionState, msg: &str| {
- check_batch_tx(batch_id, status, msg, status).await;
- };
- let check_order_tx = async |order_id: &str,
- status: SubmissionState,
- msg: &str,
- tx_status: SubmissionState| {
- let batch_id = sqlx::query(
- "SELECT initiated_outgoing_batch_id FROM initiated_outgoing_batches WHERE order_id=$1"
- ).bind(order_id)
- .try_map(|r: PgRow| {
- r.try_get_u64(0)
- }).fetch_one(&db).await.unwrap();
- check_batch_tx(batch_id, status, msg, tx_status).await;
- };
- let check_order = async |order_id: &str, status: SubmissionState, msg: &str| {
- check_order_tx(order_id, status, msg, status).await;
- };
-
- async fn test<F>(db: &PgPool, lambda: impl FnOnce(u64) -> F)
- where
- F: Future<Output = ()>,
- {
- // Reset DB
- sqlx::query("DELETE FROM initiated_outgoing_transactions")
- .execute(db)
- .await
- .unwrap();
- sqlx::query("DELETE FROM initiated_outgoing_batches")
- .execute(db)
- .await
- .unwrap();
- // Create a test batch with three transactions
- for id in ["TX", "TX_SETTLED"] {
- assert!(matches!(
- initiate(db, &gen_init_pay(id, "lol")).await.unwrap(),
- PaymentInitiationResult::Success(_)
- ));
- }
- batch_initiated(db, &Timestamp::now(), "BATCH", false)
- .await
- .unwrap();
-
- // Create witness transactions and batch
- for id in ["WITNESS_1", "WITNESS_2"] {
- assert!(matches!(
- initiate(db, &gen_init_pay(id, "lol")).await.unwrap(),
- PaymentInitiationResult::Success(_)
- ));
- }
- batch_initiated(db, &Timestamp::now(), "BATCH_WITNESS", false)
- .await
- .unwrap();
- for id in ["WITNESS_3", "WITNESS_4"] {
- assert!(matches!(
- initiate(db, &gen_init_pay(id, "lol")).await.unwrap(),
- PaymentInitiationResult::Success(_)
- ));
- }
- // Check everything is unsubmitted
- sqlx::query(
- "
- SELECT (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_batches)
- AND (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_transactions)
- "
- ).try_map(|r: PgRow| {
- assert!(r.try_get_flag(0).unwrap());
- Ok(())
- }).fetch_one(db).await.unwrap();
- let submitibale = initiated_submittable(db, &CURRENCY).await.unwrap();
- lambda(
- submitibale
- .iter()
- .find(|it| it.msg_id == "BATCH")
- .unwrap()
- .id,
- );
- // Check witness status is unaltered
- 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')
- FROM initiated_outgoing_transactions JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
- WHERE message_id != 'BATCH')
- "
- ).try_map(|r: PgRow| {
- assert!(r.try_get(0)?);
- Ok(())
- }).fetch_one(db).await.unwrap();
- }
-
- let now = Timestamp::now();
-
- // Submission retry status
- test(&db, async |batch_id| {
- batch_sub_failure(&db, batch_id, &now, "First failure")
- .await
- .unwrap();
- check_batch(batch_id, transient_failure, "First failure").await;
- batch_sub_failure(&db, batch_id, &now, "Second failure")
- .await
- .unwrap();
- check_batch(batch_id, transient_failure, "Second failure").await;
- batch_sub_success(&db, batch_id, &now, "ORDER")
- .await
- .unwrap();
- check_order("ORDER", pending, "").await;
- batch_sub_success(&db, batch_id, &now, "ORDER")
- .await
- .unwrap();
- check_order("ORDER", pending, "").await;
- order_step(&db, "ORDER", "step msg").await.unwrap();
- check_order("ORDER", pending, "step msg").await;
- order_step(&db, "ORDER", "success msg").await.unwrap();
- check_order("ORDER", pending, "success msg").await;
- order_success(&db, "ORDER").await.unwrap();
- check_order_tx("ORDER", success, "success msg", pending).await;
- order_step(&db, "ORDER", "late msg").await.unwrap();
- check_order_tx("ORDER", success, "success msg", pending).await;
- })
- .await;
-
- // Order step message on failure
- test(&db, async |batch_id| {
- batch_sub_success(&db, batch_id, &now, "ORDER")
- .await
- .unwrap();
- check_order("ORDER", pending, "").await;
- order_step(&db, "ORDER", "step msg").await.unwrap();
- check_order("ORDER", pending, "step msg").await;
- order_step(&db, "ORDER", "failure msg").await.unwrap();
- check_order("ORDER", pending, "failure msg").await;
- assert_eq!(
- Some("failure msg"),
- order_failure(&db, "ORDER")
- .await
- .unwrap()
- .unwrap()
- .1
- .as_deref()
- );
- check_order("ORDER", permanent_failure, "late msg").await;
- order_step(&db, "ORDER", "late msg").await.unwrap();
- check_order("ORDER", permanent_failure, "failure msg").await;
- })
- .await;
-
- // Payment & batch status
- test(&db, async |batch_id| {
- check_batch(batch_id, unsubmitted, "").await;
- batch_status_update(&db, "BATCH", pending, "progress")
- .await
- .unwrap();
- check_batch(batch_id, pending, "progress").await;
- tx_status_update(&db, "TX_SETTLED", "", success, "success")
- .await
- .unwrap();
- check_parts(
- batch_id, pending, "progress", pending, "progress", success, "success",
- )
- .await;
- batch_status_update(&db, "BATCH", transient_failure, "waiting")
- .await
- .unwrap();
- check_parts(
- batch_id,
- transient_failure,
- "waiting",
- transient_failure,
- "waiting",
- success,
- "success",
- )
- .await;
- tx_status_update(&db, "TX", "BATCH", permanent_failure, "failure")
- .await
- .unwrap();
- check_parts(
- batch_id,
- success,
- "",
- permanent_failure,
- "failure",
- success,
- "success",
- )
- .await;
- tx_status_update(&db, "TX_SETTLED", "BATCH", permanent_failure, "late")
- .await
- .unwrap();
- check_parts(
- batch_id,
- success,
- "",
- permanent_failure,
- "failure",
- late_failure,
- "late",
- )
- .await;
- })
- .await;
-
- // Registration
- test(&db, async |batch_id| {
- check_batch(batch_id, unsubmitted, "").await;
- register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX_SETTLED"))
- .await
- .unwrap();
- check_parts(batch_id, unsubmitted, "", unsubmitted, "", success, "").await;
- register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX").with_msg_id("BATCH"))
- .await
- .unwrap();
- check_parts(batch_id, success, "", success, "", success, "").await;
- })
- .await;
-
- // Transaction failure take over batch failures
- test(&db, async |batch_id| {
- check_batch(batch_id, unsubmitted, "").await;
- batch_status_update(&db, "BATCH", permanent_failure, "batch")
- .await
- .unwrap();
- check_parts(
- batch_id,
- permanent_failure,
- "batch",
- permanent_failure,
- "batch",
- permanent_failure,
- "batch",
- )
- .await;
- tx_status_update(&db, "TX", "BATCH", permanent_failure, "tx")
- .await
- .unwrap();
- batch_status_update(&db, "BATCH", permanent_failure, "batch2")
- .await
- .unwrap();
- check_parts(
- batch_id,
- permanent_failure,
- "batch",
- permanent_failure,
- "tx",
- permanent_failure,
- "batch",
- )
- .await;
- })
- .await;
-
- // Unknown order and batch
- batch_sub_success(&db, 42, &now, "ORDER_X").await.unwrap();
- batch_sub_failure(&db, 42, &now, "").await.unwrap();
- order_step(&db, "ORDER_X", "msg").await.unwrap();
- batch_status_update(&db, "BATCH_X", success, "")
- .await
- .unwrap();
- tx_status_update(&db, "TX_X", "BATCH_X", success, "msg")
- .await
- .unwrap();
- assert!(order_success(&db, "ORDER_X").await.unwrap().is_none());
- assert!(order_failure(&db, "ORDER_X").await.unwrap().is_none());
- }
-
- #[tokio::test]
- pub async fn initiated_submittables() {
- let (_, db) = setup().await;
- let now = Timestamp::now();
- for i in 0..6 {
- assert!(matches!(
- initiate(&db, &gen_init_pay(format!("PAY{i}"), "")).await,
- Ok(PaymentInitiationResult::Success(_))
- ));
- batch_initiated(&db, &now, &rand_ebics_id(), false)
- .await
- .unwrap();
- }
-
- let check_ids = async |ids: &[&str]| {
- assert_eq!(
- ids,
- initiated_submittable(&db, &CURRENCY)
- .await
- .unwrap()
- .iter()
- .flat_map(|it| it.payments.iter().map(|it| it.end_to_end_id.as_str()))
- .collect::<Vec<_>>()
- );
- };
- check_ids(&["PAY0", "PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await;
-
- // Check submitted not submitable
- batch_sub_success(&db, 1, &now, "ORDER1").await.unwrap();
- check_ids(&["PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await;
-
- // Check transient failure submitable last
- batch_sub_failure(&db, 2, &now, "Failure").await.unwrap();
- check_ids(&["PAY2", "PAY3", "PAY4", "PAY5", "PAY1"]).await;
-
- // Check persistent failure not submitable
- batch_sub_success(&db, 4, &now, "ORDER3").await.unwrap();
- order_failure(&db, "ORDER3").await.unwrap();
- check_ids(&["PAY2", "PAY4", "PAY5", "PAY1"]).await;
- batch_sub_success(&db, 5, &now, "ORDER4").await.unwrap();
- order_failure(&db, "ORDER4").await.unwrap();
- check_ids(&["PAY2", "PAY5", "PAY1"]).await;
-
- // Check rotation
- batch_sub_failure(&db, 3, &Timestamp::now(), "FAILURE")
- .await
- .unwrap();
- check_ids(&["PAY5", "PAY1", "PAY2"]).await;
- batch_sub_failure(&db, 6, &Timestamp::now(), "FAILURE")
- .await
- .unwrap();
- check_ids(&["PAY1", "PAY2", "PAY5"]).await;
- batch_sub_failure(&db, 2, &Timestamp::now(), "FAILURE")
- .await
- .unwrap();
- check_ids(&["PAY2", "PAY5", "PAY1"]).await;
- }
}
diff --git a/src/db/initiated.rs b/src/db/initiated.rs
@@ -0,0 +1,886 @@
+/*
+ This file is part of TALER
+ Copyright (C) 2026 Taler Systems SA
+
+ TALER is free software; you can redistribute it and/or modify it under the
+ terms of the GNU Affero General Public License as published by the Free Software
+ Foundation; either version 3, or (at your option) any later version.
+
+ TALER is distributed in the hope that it will be useful, but WITHOUT ANY
+ WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR
+ A PARTICULAR PURPOSE. See the GNU Affero General Public License for more details.
+
+ You should have received a copy of the GNU Affero General Public License along with
+ TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/>
+*/
+
+use std::collections::BTreeMap;
+
+use const_format::formatcp;
+use jiff::Timestamp;
+use sqlx::{PgPool, Row as _, postgres::PgRow};
+use taler_api::db::{BindHelper as _, TypeHelper as _};
+use taler_common::types::amount::Currency;
+
+use crate::{
+ db::{PENDING, UNSETTLED},
+ model::{InitiatedPayment, OutgoingId, OutgoingPayment, PaymentBatch, SubmissionState},
+};
+
+/// Outgoing payments initiation result
+#[derive(Debug, PartialEq, Eq)]
+pub enum PaymentInitiationResult {
+ Success(u64),
+ RequestUidReuse,
+}
+
+/// Initiate a new payment
+pub async fn initiate(
+ pool: &PgPool,
+ payment: &InitiatedPayment,
+) -> sqlx::Result<PaymentInitiationResult> {
+ let res = sqlx::query(
+ "
+ INSERT INTO initiated_outgoing_transactions (
+ amount,
+ subject,
+ credit_payto,
+ initiation_time,
+ end_to_end_id
+ ) VALUES ($1,$2,$3,$4,$5)
+ RETURNING initiated_outgoing_transaction_id
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(&payment.subject)
+ .bind(payment.creditor.as_ref().as_str())
+ .bind_timestamp(&payment.initiation_time)
+ .bind(&payment.end_to_end_id)
+ .try_map(|r: PgRow| Ok(PaymentInitiationResult::Success(r.try_get_u64(0)?)))
+ .fetch_one(pool)
+ .await;
+ if let Err(e) = &res
+ && let Some(db_err) = e.as_database_error()
+ && db_err.code() == Some(std::borrow::Cow::Borrowed("23505"))
+ {
+ Ok(PaymentInitiationResult::RequestUidReuse)
+ } else {
+ res
+ }
+}
+
+/// Group unbatched transaction into a single batch
+pub async fn batch_initiated(
+ pool: &PgPool,
+ timestamp: &Timestamp,
+ ebics_id: &str,
+ require_ack: bool,
+) -> sqlx::Result<()> {
+ sqlx::query("SELECT batch_outgoing_transactions($1, $2, $3)")
+ .bind_timestamp(timestamp)
+ .bind(ebics_id)
+ .bind(require_ack)
+ .execute(pool)
+ .await?;
+ Ok(())
+}
+
+pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<()> {
+ sqlx::query("UPDATE initiated_outgoing_transactions SET awaiting_ack=false WHERE initiated_outgoing_transaction_id=$1")
+ .bind(id as i64)
+ .execute(db)
+ .await?;
+ Ok(())
+}
+
+pub async fn initiated_submittable(
+ db: &PgPool,
+ currency: &Currency,
+) -> sqlx::Result<Vec<PaymentBatch>> {
+ const SELECT_PART: &str = "
+ 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!(
+ "
+ ({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(),
+ })
+ })
+ .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
+ ,subject
+ ,credit_payto
+ ,initiated_outgoing_transactions.initiation_time
+ ,end_to_end_id
+ ,initiated_outgoing_batch_id
+ FROM initiated_outgoing_transactions
+ JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
+ WHERE initiated_outgoing_batches.status IN ('unsubmitted', 'transient_failure')
+ ",
+ )
+ .try_map(|r: PgRow| {
+ let payment = InitiatedPayment {
+ 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")?,
+ end_to_end_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)
+}
+
+pub async fn unsettled_tx_in_batch(
+ db: &PgPool,
+ currency: &Currency,
+ msg_id: &str,
+ execution_time: &Timestamp,
+) -> sqlx::Result<Vec<OutgoingPayment>> {
+ sqlx::query(formatcp!(
+ "
+ SELECT
+ end_to_end_id,
+ amount,
+ subject,
+ credit_payto
+ FROM initiated_outgoing_transactions
+ JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
+ WHERE message_id = $1
+ AND initiated_outgoing_transactions.{UNSETTLED}
+ "
+ ))
+ .bind(msg_id)
+ .try_map(|r: PgRow| {
+ Ok(OutgoingPayment {
+ id: OutgoingId {
+ msg_id: Some(msg_id.into()),
+ end_to_end_id: r.try_get("end_to_end_id")?,
+ acct_svcr_ref: None,
+ },
+ amount: r.try_get_amount("amount", currency)?,
+ debit_fee: None,
+ subject: r.try_get("subject")?,
+ execution_time: *execution_time,
+ creditor: r.try_get_opt_payto("credit_payto")?,
+ })
+ })
+ .fetch_all(db)
+ .await
+}
+
+/** Register submission success of order [orderId] for batch [id] at [timestamp] */
+pub async fn batch_sub_success(
+ db: &PgPool,
+ batch_id: u64,
+ timestamp: &Timestamp,
+ order_id: &str,
+) -> sqlx::Result<()> {
+ 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
+ ,order_id = $2
+ ,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(batch_id as i64)
+ .execute(&mut *tx)
+ .await?;
+ }
+ tx.commit().await
+}
+
+/** Register submission failure with [msg] for batch [id] at [timestamp]*/
+pub async fn batch_sub_failure(
+ db: &PgPool,
+ batch_id: u64,
+ timestamp: &Timestamp,
+ msg: &str,
+) -> sqlx::Result<()> {
+ let permanent = false;
+ let mut tx = db.begin().await?;
+ // Update batch status
+ sqlx::query(
+ "
+ UPDATE initiated_outgoing_batches
+ SET status = $1
+ ,submission_date = $2
+ ,status_msg = $3
+ ,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
+ 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(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
+ SET status = 'pending', status_msg = $1
+ WHERE order_id = $1 AND {PENDING}
+ RETURNING initiated_outgoing_batch_id
+ "
+ ))
+ .bind(msg)
+ .bind(order_id)
+ .try_map(|r: PgRow| r.try_get_u32(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)
+ .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
+ 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_u32(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))
+}
+
+/** Register order failure for [orderId] and return message_id and previous status_msg if found */
+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
+ 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
+ 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)))
+}
+
+/** Register payment status [state] with [msg] for batch [msgId] */
+pub async fn batch_status_update(
+ db: &PgPool,
+ msg_id: &str,
+ 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
+}
+
+/** Register payment status [state] with [msg] for transaction [endToEndId] in batch [msgId] */
+pub async fn tx_status_update(
+ db: &PgPool,
+ end_to_end_id: &str,
+ msg_id: &str,
+ 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
+}
+
+#[cfg(test)]
+mod test {
+ use std::str::FromStr as _;
+
+ use jiff::{Span, Timestamp, civil::Date};
+ use sqlx::{PgPool, Row as _, postgres::PgRow};
+ use taler_api::db::TypeHelper as _;
+ use taler_common::{config::Config, types::utils::date_to_utc_ts};
+
+ use crate::{
+ CONFIG_SOURCE,
+ config::{NexusCfg, NexusIngestConfig},
+ db::{
+ initiated::{
+ PaymentInitiationResult, batch_initiated, batch_status_update, batch_sub_failure,
+ batch_sub_success, initiate, initiated_submittable, order_failure, order_step,
+ order_success, tx_status_update,
+ },
+ test::{CURRENCY, check_count, gen_in_pay, gen_init_pay, gen_out_pay, setup},
+ },
+ model::{SubmissionState, Tx},
+ rand_ebics_id,
+ worker::{register_outgoing, register_tx},
+ };
+
+ #[tokio::test]
+ pub async fn initiated_skip() {
+ let (_, db) = setup().await;
+ let cfg = Config::from_file(CONFIG_SOURCE, Some("libeufin-nexus/conf/skip.conf")).unwrap();
+ let cfg = NexusCfg::parse(cfg).unwrap();
+ let cfg = cfg.ingest().unwrap();
+ let millis = Span::new().milliseconds(10);
+
+ async fn ingest(db: &PgPool, cfg: &NexusIngestConfig, execution_time: Timestamp) {
+ for tx in [
+ Tx::In(
+ gen_in_pay(format!("test at {execution_time}"))
+ .with_execution_time(execution_time),
+ ),
+ Tx::Out(
+ gen_out_pay(format!("test at {execution_time}"))
+ .with_execution_time(execution_time),
+ ),
+ ] {
+ register_tx(db, cfg, &tx).await.unwrap()
+ }
+ }
+
+ assert_eq!(
+ cfg.ignore_txs_before,
+ date_to_utc_ts(&Date::from_str("2024-04-04").unwrap())
+ );
+ assert_eq!(
+ cfg.ignore_bounces_before,
+ date_to_utc_ts(&Date::from_str("2024-06-12").unwrap())
+ );
+
+ // No transaction at the beginning
+ check_count(&db, 0, 0).await;
+
+ // Skipped transactions
+ ingest(&db, &cfg, cfg.ignore_txs_before - millis).await;
+ check_count(&db, 0, 0).await;
+
+ // Skipped bounces
+ ingest(&db, &cfg, cfg.ignore_txs_before).await;
+ ingest(&db, &cfg, cfg.ignore_txs_before + millis).await;
+ ingest(&db, &cfg, cfg.ignore_bounces_before - millis).await;
+ check_count(&db, 6, 0).await;
+
+ // Bounces
+ ingest(&db, &cfg, cfg.ignore_bounces_before).await;
+ ingest(&db, &cfg, cfg.ignore_bounces_before + millis).await;
+ check_count(&db, 10, 2).await;
+ }
+
+ #[tokio::test]
+ pub async fn initiated_status() {
+ use SubmissionState::*;
+
+ let (_, db) = setup().await;
+
+ let check_parts = async |batch_id: u64,
+ batch_status: SubmissionState,
+ batch_msg: &str,
+ tx_status: SubmissionState,
+ tx_msg: &str,
+ settled_status: SubmissionState,
+ settled_msg: &str| {
+ // Check batch status
+ let msg_id: String = sqlx::query(
+ "
+ SELECT message_id, status, status_msg FROM initiated_outgoing_batches WHERE initiated_outgoing_batch_id=?
+ "
+ ).bind(batch_id as i64)
+ .try_map(|r: PgRow| {
+ let msg_id: String = r.try_get("message_id")?;
+ assert_eq!((batch_status, batch_msg), (r.try_get("status")?, r.try_get("status_msg")?), "{msg_id}");
+ Ok(msg_id)
+ }).fetch_one(&db).await.unwrap();
+ // Check tx status
+ sqlx::query(
+ "
+ SELECT end_to_end_id, status, status_msg FROM initiated_outgoing_transactions WHERE initiated_outgoing_batch_id=?
+ "
+ ).bind(batch_id as i64).try_map(|r: PgRow| {
+ let end_to_end_id: &str = r.try_get("end_to_end_id")?;
+ let expected = match end_to_end_id {
+ "TX" => (tx_status, tx_msg),
+ "TX_SETTLED" => (settled_status, settled_msg),
+ _ =>panic!("Unexpected tx $endToEndId")
+ };
+ assert_eq!(expected,
+ (r.try_get("status")?, r.try_get("status_msg")?),
+ "{msg_id},{end_to_end_id}"
+ );
+ Ok(())
+ }).fetch_all(&db).await.unwrap();
+ };
+
+ let check_batch_tx = async |batch_id: u64,
+ status: SubmissionState,
+ msg: &str,
+ tx_status: SubmissionState| {
+ check_parts(batch_id, status, msg, tx_status, msg, tx_status, msg).await;
+ };
+ let check_batch = async |batch_id: u64, status: SubmissionState, msg: &str| {
+ check_batch_tx(batch_id, status, msg, status).await;
+ };
+ let check_order_tx = async |order_id: &str,
+ status: SubmissionState,
+ msg: &str,
+ tx_status: SubmissionState| {
+ let batch_id = sqlx::query(
+ "SELECT initiated_outgoing_batch_id FROM initiated_outgoing_batches WHERE order_id=$1"
+ ).bind(order_id)
+ .try_map(|r: PgRow| {
+ r.try_get_u64(0)
+ }).fetch_one(&db).await.unwrap();
+ check_batch_tx(batch_id, status, msg, tx_status).await;
+ };
+ let check_order = async |order_id: &str, status: SubmissionState, msg: &str| {
+ check_order_tx(order_id, status, msg, status).await;
+ };
+
+ async fn test<F>(db: &PgPool, lambda: impl FnOnce(u64) -> F)
+ where
+ F: Future<Output = ()>,
+ {
+ // Reset DB
+ sqlx::query("DELETE FROM initiated_outgoing_transactions")
+ .execute(db)
+ .await
+ .unwrap();
+ sqlx::query("DELETE FROM initiated_outgoing_batches")
+ .execute(db)
+ .await
+ .unwrap();
+ // Create a test batch with three transactions
+ for id in ["TX", "TX_SETTLED"] {
+ assert!(matches!(
+ initiate(db, &gen_init_pay(id, "lol")).await.unwrap(),
+ PaymentInitiationResult::Success(_)
+ ));
+ }
+ batch_initiated(db, &Timestamp::now(), "BATCH", false)
+ .await
+ .unwrap();
+
+ // Create witness transactions and batch
+ for id in ["WITNESS_1", "WITNESS_2"] {
+ assert!(matches!(
+ initiate(db, &gen_init_pay(id, "lol")).await.unwrap(),
+ PaymentInitiationResult::Success(_)
+ ));
+ }
+ batch_initiated(db, &Timestamp::now(), "BATCH_WITNESS", false)
+ .await
+ .unwrap();
+ for id in ["WITNESS_3", "WITNESS_4"] {
+ assert!(matches!(
+ initiate(db, &gen_init_pay(id, "lol")).await.unwrap(),
+ PaymentInitiationResult::Success(_)
+ ));
+ }
+ // Check everything is unsubmitted
+ sqlx::query(
+ "
+ SELECT (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_batches)
+ AND (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_transactions)
+ "
+ ).try_map(|r: PgRow| {
+ assert!(r.try_get_flag(0).unwrap());
+ Ok(())
+ }).fetch_one(db).await.unwrap();
+ let submitibale = initiated_submittable(db, &CURRENCY).await.unwrap();
+ lambda(
+ submitibale
+ .iter()
+ .find(|it| it.msg_id == "BATCH")
+ .unwrap()
+ .id,
+ );
+ // Check witness status is unaltered
+ 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')
+ FROM initiated_outgoing_transactions JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
+ WHERE message_id != 'BATCH')
+ "
+ ).try_map(|r: PgRow| {
+ assert!(r.try_get(0)?);
+ Ok(())
+ }).fetch_one(db).await.unwrap();
+ }
+
+ let now = Timestamp::now();
+
+ // Submission retry status
+ test(&db, async |batch_id| {
+ batch_sub_failure(&db, batch_id, &now, "First failure")
+ .await
+ .unwrap();
+ check_batch(batch_id, transient_failure, "First failure").await;
+ batch_sub_failure(&db, batch_id, &now, "Second failure")
+ .await
+ .unwrap();
+ check_batch(batch_id, transient_failure, "Second failure").await;
+ batch_sub_success(&db, batch_id, &now, "ORDER")
+ .await
+ .unwrap();
+ check_order("ORDER", pending, "").await;
+ batch_sub_success(&db, batch_id, &now, "ORDER")
+ .await
+ .unwrap();
+ check_order("ORDER", pending, "").await;
+ order_step(&db, "ORDER", "step msg").await.unwrap();
+ check_order("ORDER", pending, "step msg").await;
+ order_step(&db, "ORDER", "success msg").await.unwrap();
+ check_order("ORDER", pending, "success msg").await;
+ order_success(&db, "ORDER").await.unwrap();
+ check_order_tx("ORDER", success, "success msg", pending).await;
+ order_step(&db, "ORDER", "late msg").await.unwrap();
+ check_order_tx("ORDER", success, "success msg", pending).await;
+ })
+ .await;
+
+ // Order step message on failure
+ test(&db, async |batch_id| {
+ batch_sub_success(&db, batch_id, &now, "ORDER")
+ .await
+ .unwrap();
+ check_order("ORDER", pending, "").await;
+ order_step(&db, "ORDER", "step msg").await.unwrap();
+ check_order("ORDER", pending, "step msg").await;
+ order_step(&db, "ORDER", "failure msg").await.unwrap();
+ check_order("ORDER", pending, "failure msg").await;
+ assert_eq!(
+ Some("failure msg"),
+ order_failure(&db, "ORDER")
+ .await
+ .unwrap()
+ .unwrap()
+ .1
+ .as_deref()
+ );
+ check_order("ORDER", permanent_failure, "late msg").await;
+ order_step(&db, "ORDER", "late msg").await.unwrap();
+ check_order("ORDER", permanent_failure, "failure msg").await;
+ })
+ .await;
+
+ // Payment & batch status
+ test(&db, async |batch_id| {
+ check_batch(batch_id, unsubmitted, "").await;
+ batch_status_update(&db, "BATCH", pending, "progress")
+ .await
+ .unwrap();
+ check_batch(batch_id, pending, "progress").await;
+ tx_status_update(&db, "TX_SETTLED", "", success, "success")
+ .await
+ .unwrap();
+ check_parts(
+ batch_id, pending, "progress", pending, "progress", success, "success",
+ )
+ .await;
+ batch_status_update(&db, "BATCH", transient_failure, "waiting")
+ .await
+ .unwrap();
+ check_parts(
+ batch_id,
+ transient_failure,
+ "waiting",
+ transient_failure,
+ "waiting",
+ success,
+ "success",
+ )
+ .await;
+ tx_status_update(&db, "TX", "BATCH", permanent_failure, "failure")
+ .await
+ .unwrap();
+ check_parts(
+ batch_id,
+ success,
+ "",
+ permanent_failure,
+ "failure",
+ success,
+ "success",
+ )
+ .await;
+ tx_status_update(&db, "TX_SETTLED", "BATCH", permanent_failure, "late")
+ .await
+ .unwrap();
+ check_parts(
+ batch_id,
+ success,
+ "",
+ permanent_failure,
+ "failure",
+ late_failure,
+ "late",
+ )
+ .await;
+ })
+ .await;
+
+ // Registration
+ test(&db, async |batch_id| {
+ check_batch(batch_id, unsubmitted, "").await;
+ register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX_SETTLED"))
+ .await
+ .unwrap();
+ check_parts(batch_id, unsubmitted, "", unsubmitted, "", success, "").await;
+ register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX").with_msg_id("BATCH"))
+ .await
+ .unwrap();
+ check_parts(batch_id, success, "", success, "", success, "").await;
+ })
+ .await;
+
+ // Transaction failure take over batch failures
+ test(&db, async |batch_id| {
+ check_batch(batch_id, unsubmitted, "").await;
+ batch_status_update(&db, "BATCH", permanent_failure, "batch")
+ .await
+ .unwrap();
+ check_parts(
+ batch_id,
+ permanent_failure,
+ "batch",
+ permanent_failure,
+ "batch",
+ permanent_failure,
+ "batch",
+ )
+ .await;
+ tx_status_update(&db, "TX", "BATCH", permanent_failure, "tx")
+ .await
+ .unwrap();
+ batch_status_update(&db, "BATCH", permanent_failure, "batch2")
+ .await
+ .unwrap();
+ check_parts(
+ batch_id,
+ permanent_failure,
+ "batch",
+ permanent_failure,
+ "tx",
+ permanent_failure,
+ "batch",
+ )
+ .await;
+ })
+ .await;
+
+ // Unknown order and batch
+ batch_sub_success(&db, 42, &now, "ORDER_X").await.unwrap();
+ batch_sub_failure(&db, 42, &now, "").await.unwrap();
+ order_step(&db, "ORDER_X", "msg").await.unwrap();
+ batch_status_update(&db, "BATCH_X", success, "")
+ .await
+ .unwrap();
+ tx_status_update(&db, "TX_X", "BATCH_X", success, "msg")
+ .await
+ .unwrap();
+ assert!(order_success(&db, "ORDER_X").await.unwrap().is_none());
+ assert!(order_failure(&db, "ORDER_X").await.unwrap().is_none());
+ }
+
+ #[tokio::test]
+ pub async fn initiated_submittables() {
+ let (_, db) = setup().await;
+ let now = Timestamp::now();
+ for i in 0..6 {
+ assert!(matches!(
+ initiate(&db, &gen_init_pay(format!("PAY{i}"), "")).await,
+ Ok(PaymentInitiationResult::Success(_))
+ ));
+ batch_initiated(&db, &now, &rand_ebics_id(), false)
+ .await
+ .unwrap();
+ }
+
+ let check_ids = async |ids: &[&str]| {
+ assert_eq!(
+ ids,
+ initiated_submittable(&db, &CURRENCY)
+ .await
+ .unwrap()
+ .iter()
+ .flat_map(|it| it.payments.iter().map(|it| it.end_to_end_id.as_str()))
+ .collect::<Vec<_>>()
+ );
+ };
+ check_ids(&["PAY0", "PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await;
+
+ // Check submitted not submitable
+ batch_sub_success(&db, 1, &now, "ORDER1").await.unwrap();
+ check_ids(&["PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await;
+
+ // Check transient failure submitable last
+ batch_sub_failure(&db, 2, &now, "Failure").await.unwrap();
+ check_ids(&["PAY2", "PAY3", "PAY4", "PAY5", "PAY1"]).await;
+
+ // Check persistent failure not submitable
+ batch_sub_success(&db, 4, &now, "ORDER3").await.unwrap();
+ order_failure(&db, "ORDER3").await.unwrap();
+ check_ids(&["PAY2", "PAY4", "PAY5", "PAY1"]).await;
+ batch_sub_success(&db, 5, &now, "ORDER4").await.unwrap();
+ order_failure(&db, "ORDER4").await.unwrap();
+ check_ids(&["PAY2", "PAY5", "PAY1"]).await;
+
+ // Check rotation
+ batch_sub_failure(&db, 3, &Timestamp::now(), "FAILURE")
+ .await
+ .unwrap();
+ check_ids(&["PAY5", "PAY1", "PAY2"]).await;
+ batch_sub_failure(&db, 6, &Timestamp::now(), "FAILURE")
+ .await
+ .unwrap();
+ check_ids(&["PAY1", "PAY2", "PAY5"]).await;
+ batch_sub_failure(&db, 2, &Timestamp::now(), "FAILURE")
+ .await
+ .unwrap();
+ check_ids(&["PAY2", "PAY5", "PAY1"]).await;
+ }
+}
diff --git a/src/db/payment.rs b/src/db/payment.rs
@@ -0,0 +1,1169 @@
+/*
+ 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 compact_str::CompactString;
+use jiff::Timestamp;
+use sqlx::{PgPool, Row as _, postgres::PgRow};
+use taler_api::{
+ db::{BindHelper as _, TypeHelper as _},
+ subject::{IncomingSubject, OutgoingSubject},
+};
+use taler_common::types::amount::{Amount, Currency};
+
+use crate::model::{IncomingPayment, OutgoingPayment};
+
+#[derive(Debug, PartialEq, Eq)]
+pub struct OutgoingRegistrationResult {
+ pub id: u64,
+ pub initiated: bool,
+ pub new: bool,
+}
+
+/** Register an outgoing payment reconciling it with its initiated payment counterpart if present */
+pub async fn register_out_tx(
+ pool: &PgPool,
+ payment: &OutgoingPayment,
+ subject: Option<&OutgoingSubject>,
+) -> sqlx::Result<OutgoingRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT out_tx_id, out_initiated, out_found
+ FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .debit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(&payment.id.end_to_end_id)
+ .bind(&payment.id.msg_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(subject.as_ref().map(|s| &s.wtid))
+ .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str()))
+ .bind(subject.as_ref().map(|s| &s.metadata))
+ .try_map(|r: PgRow| {
+ Ok(OutgoingRegistrationResult {
+ id: r.try_get_u64(0)?,
+ initiated: r.try_get_flag(1)?,
+ new: !r.try_get_flag(2)?,
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+/// Register an outgoing batch
+pub async fn register_out_batch(
+ pool: &PgPool,
+ currency: &Currency,
+ payment: &OutgoingPayment,
+ subject: Option<&OutgoingSubject>,
+) -> sqlx::Result<OutgoingRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT out_tx_id, out_initiated, out_found
+ FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .debit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(&payment.id.end_to_end_id)
+ .bind(&payment.id.msg_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(subject.as_ref().map(|s| &s.wtid))
+ .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str()))
+ .bind(subject.as_ref().map(|s| &s.metadata))
+ .try_map(|r: PgRow| {
+ Ok(OutgoingRegistrationResult {
+ id: r.try_get_u64(0)?,
+ initiated: r.try_get_flag(1)?,
+ new: !r.try_get_flag(2)?,
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub struct InResult {
+ pub id: u64,
+ pub new: bool,
+ pub completed: bool,
+ pub pending: bool,
+ pub bounce_id: Option<CompactString>,
+}
+
+/** Incoming payments registration result */
+#[derive(Debug, PartialEq, Eq)]
+pub enum IncomingRegistrationResult {
+ Success(InResult),
+ ReservePubReuse,
+ MappingReuse,
+ UnknownMapping,
+}
+
+/** Register an incoming payment */
+pub async fn register_in(pool: &PgPool, payment: &IncomingPayment) -> sqlx::Result<InResult> {
+ sqlx::query(
+ "
+ SELECT out_found, out_completed, out_tx_id, out_bounce_id
+ FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,NULL)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .try_map(|r: PgRow| {
+ Ok(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: false,
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+/** Register an talerable incoming payment */
+pub async fn register_in_talerable(
+ pool: &PgPool,
+ payment: &IncomingPayment,
+ subject: &IncomingSubject,
+) -> sqlx::Result<IncomingRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT
+ out_reserve_pub_reuse,
+ out_mapping_reuse,
+ out_unknown_mapping,
+ out_found,
+ out_completed,
+ out_pending,
+ out_tx_id,
+ out_bounce_id
+ FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9::taler_incoming_type,$10,NULL)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(subject.ty())
+ .bind(subject.key())
+ .try_map(|r: PgRow| {
+ Ok(if r.try_get_flag("out_reserve_pub_reuse")? {
+ IncomingRegistrationResult::ReservePubReuse
+ } else if r.try_get_flag("out_mapping_reuse")? {
+ IncomingRegistrationResult::MappingReuse
+ } else if r.try_get_flag("out_unknown_mapping")? {
+ IncomingRegistrationResult::UnknownMapping
+ } else {
+ IncomingRegistrationResult::Success(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: r.try_get("out_pending")?,
+ })
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+/** Register an talerable incoming payment */
+pub async fn register_in_qr_bill(
+ pool: &PgPool,
+ payment: &IncomingPayment,
+ reference: &str,
+) -> sqlx::Result<IncomingRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT
+ out_reserve_pub_reuse,
+ out_mapping_reuse,
+ out_unknown_mapping,
+ out_found,
+ out_completed,
+ out_pending,
+ out_tx_id,
+ out_bounce_id
+ FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,$9)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(reference)
+ .try_map(|r: PgRow| {
+ Ok(if r.try_get_flag("out_reserve_pub_reuse")? {
+ IncomingRegistrationResult::ReservePubReuse
+ } else if r.try_get_flag("out_mapping_reuse")? {
+ IncomingRegistrationResult::MappingReuse
+ } else if r.try_get_flag("out_unknown_mapping")? {
+ IncomingRegistrationResult::UnknownMapping
+ } else {
+ IncomingRegistrationResult::Success(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: r.try_get("out_pending")?,
+ })
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+/** Incoming payments bounce registration result */
+pub enum IncomingBounceRegistrationResult {
+ Success(InResult),
+ Talerable,
+}
+
+/** Register an incoming payment and bounce it */
+pub async fn register_in_malformed(
+ pool: &PgPool,
+ payment: &IncomingPayment,
+ bounce_amount: &Amount,
+ bounce_end_to_end_id: &str,
+ timestamp: &Timestamp,
+ cause: &str,
+) -> sqlx::Result<IncomingBounceRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT out_found, out_tx_id, out_completed, out_bounce_id, out_talerable
+ FROM register_and_bounce_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(bounce_amount)
+ .bind_timestamp(timestamp)
+ .bind(bounce_end_to_end_id)
+ .bind(cause)
+ .try_map(|r: PgRow| {
+ Ok(if r.try_get_flag("out_talerable")? {
+ IncomingBounceRegistrationResult::Talerable
+ } else {
+ IncomingBounceRegistrationResult::Success(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: false,
+ })
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+#[cfg(test)]
+mod test {
+
+ use jiff::Timestamp;
+ use sqlx::{PgPool, postgres::PgRow};
+ use taler_api::{
+ db::{IncomingType, TypeHelper as _},
+ subject::subject_fmt_qr_bill,
+ };
+ use taler_common::{
+ api_common::{EddsaPublicKey, EddsaSignature, ShortHashCode},
+ types::amount::amount,
+ };
+ 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,
+ gen_init_pay, gen_out_pay, setup,
+ },
+ transfer_register,
+ },
+ model::{IncomingId, IncomingPayment, OutgoingBatch, OutgoingId, OutgoingPayment},
+ rand_ebics_id,
+ worker::{register_incoming, register_outgoing, register_outgoing_batch},
+ };
+
+ #[tokio::test]
+ async fn out_tx() {
+ let (_, db) = setup().await;
+ // Register initiated transactions
+ for subject in [
+ "initiated by nexus".to_owned(),
+ format!("{} https://exchange.com/", ShortHashCode::rand()),
+ ] {
+ let payment = gen_out_pay(subject.clone());
+ assert!(matches!(
+ initiate(
+ &db,
+ &gen_init_pay(payment.id.end_to_end_id.clone().unwrap(), subject),
+ )
+ .await,
+ Ok(PaymentInitiationResult::Success(_))
+ ));
+ let first = register_outgoing(&db, &payment).await.unwrap();
+ assert_eq!(
+ first,
+ OutgoingRegistrationResult {
+ id: first.id,
+ initiated: true,
+ new: true
+ }
+ );
+ assert_eq!(
+ register_outgoing(&db, &payment).await.unwrap(),
+ OutgoingRegistrationResult {
+ id: first.id,
+ initiated: true,
+ new: false
+ }
+ );
+ let payment = OutgoingPayment {
+ id: OutgoingId {
+ msg_id: None,
+ end_to_end_id: None,
+ acct_svcr_ref: payment.id.end_to_end_id,
+ },
+ ..payment
+ };
+ let second = register_outgoing(&db, &payment).await.unwrap();
+ assert_eq!(
+ second,
+ OutgoingRegistrationResult {
+ id: first.id + 1,
+ initiated: false,
+ new: true
+ }
+ );
+ assert_eq!(
+ register_outgoing(&db, &payment).await.unwrap(),
+ OutgoingRegistrationResult {
+ id: second.id,
+ initiated: false,
+ new: false
+ }
+ );
+ }
+ check_out_count(&db, 4, 1).await;
+
+ // Register unknown
+ for subject in [
+ "initiated by nexus".to_owned(),
+ format!("{} https://exchange.com/", ShortHashCode::rand()),
+ ] {
+ let payment = gen_out_pay(subject.clone());
+ let res = register_outgoing(&db, &payment).await.unwrap();
+ assert_eq!(
+ res,
+ OutgoingRegistrationResult {
+ id: res.id,
+ initiated: false,
+ new: true
+ }
+ );
+ assert_eq!(
+ register_outgoing(&db, &payment).await.unwrap(),
+ OutgoingRegistrationResult {
+ id: res.id,
+ initiated: false,
+ new: false
+ }
+ );
+ }
+ check_out_count(&db, 6, 2).await;
+
+ // Register wtid reuse
+ let wtid = ShortHashCode::rand();
+ for subject in [
+ format!("{wtid} https://exchange.com/"),
+ format!("{wtid} https://exchange.com/"),
+ ] {
+ let payment = gen_out_pay(subject.clone());
+ let res = register_outgoing(&db, &payment).await.unwrap();
+ assert_eq!(
+ res,
+ OutgoingRegistrationResult {
+ id: res.id,
+ initiated: false,
+ new: true
+ }
+ );
+ assert_eq!(
+ register_outgoing(&db, &payment).await.unwrap(),
+ OutgoingRegistrationResult {
+ id: res.id,
+ initiated: false,
+ new: false
+ }
+ );
+ }
+ check_out_count(&db, 8, 3).await
+ }
+
+ #[tokio::test]
+ async fn out_batch() {
+ let (_, db) = setup().await;
+ // Init batch
+ let wtid = ShortHashCode::rand();
+ for subject in [
+ "initiated by nexus".to_string(),
+ format!("{} https://exchange.com/", ShortHashCode::rand()),
+ format!("{wtid} https://exchange.com/"),
+ format!("{wtid} https://exchange.com/"),
+ ] {
+ assert!(matches!(
+ initiate(&db, &gen_init_pay(rand_ebics_id(), subject),).await,
+ Ok(PaymentInitiationResult::Success(_))
+ ));
+ }
+ batch_initiated(&db, &Timestamp::now(), "BATCH", false)
+ .await
+ .unwrap();
+
+ // Register batch
+ register_outgoing_batch(
+ &db,
+ &CURRENCY,
+ &OutgoingBatch {
+ msg_id: "BATCH".into(),
+ execution_time: Timestamp::now(),
+ },
+ )
+ .await
+ .unwrap();
+ check_out_count(&db, 4, 2).await;
+
+ // Test manual ack
+ let mut txs = Vec::new();
+ for nb in 0..3 {
+ let res = initiate(&db, &gen_init_pay(rand_ebics_id(), format!("tx {nb}"))).await;
+ if let Ok(PaymentInitiationResult::Success(id)) = &res {
+ txs.push(*id);
+ } else {
+ panic!("Expected success got {res:?}");
+ }
+ }
+
+ // Check not sent without ack
+ batch_initiated(&db, &Timestamp::now(), "BATCH_MANUAL", true)
+ .await
+ .unwrap();
+ register_outgoing_batch(
+ &db,
+ &CURRENCY,
+ &OutgoingBatch {
+ msg_id: "BATCH_MANUAL".into(),
+ execution_time: Timestamp::now(),
+ },
+ )
+ .await
+ .unwrap();
+ check_out_count(&db, 4, 2).await;
+
+ // Check sent with ack
+ for tx in txs {
+ initiated_ack(&db, tx).await.unwrap();
+ }
+ batch_initiated(&db, &Timestamp::now(), "BATCH_MANUAL", true)
+ .await
+ .unwrap();
+ register_outgoing_batch(
+ &db,
+ &CURRENCY,
+ &OutgoingBatch {
+ msg_id: "BATCH_MANUAL".into(),
+ execution_time: Timestamp::now(),
+ },
+ )
+ .await
+ .unwrap();
+ check_out_count(&db, 7, 2).await;
+ }
+
+ #[tokio::test]
+ async fn in_bounce() {
+ let (_, db) = setup().await;
+
+ // Creating and bouncing one incoming transaction
+ let payment = gen_in_pay("incoming and bounce");
+ let id = rand_ebics_id();
+
+ let bounce_amount = amount("KUDOS:2.53");
+ let res = register_in_malformed(
+ &db,
+ &payment,
+ &bounce_amount,
+ &id,
+ &Timestamp::now(),
+ "manual bounce",
+ )
+ .await
+ .unwrap();
+ assert!(
+ matches!(
+ res,
+ IncomingBounceRegistrationResult::Success(InResult {
+ new: true,
+ id: _,
+ completed: false,
+ pending: false,
+ ref bounce_id
+ }) if bounce_id.as_ref() == Some(&id)
+ ),
+ "{res:?}"
+ );
+ // Idempotent
+ let res = register_in_malformed(
+ &db,
+ &payment,
+ &amount("KUDOS:2.5"),
+ &rand_ebics_id(),
+ &Timestamp::now(),
+ "other reason to bounce",
+ )
+ .await
+ .unwrap();
+ assert!(
+ matches!(
+ res,
+ IncomingBounceRegistrationResult::Success(InResult {
+ new: false,
+ id: _,
+ completed: false,
+ pending: false,
+ ref bounce_id
+ }) if bounce_id.as_ref() == Some(&id)
+ ),
+ "{res:?}"
+ );
+
+ // Checking one incoming got created and bounced
+ sqlx::query(
+ "
+ SELECT
+ incoming_transactions.amount as in_amount,
+ initiated_outgoing_transactions.amount as bounce_amount
+ FROM incoming_transactions
+ JOIN bounced_transactions USING (incoming_transaction_id)
+ JOIN initiated_outgoing_transactions USING (initiated_outgoing_transaction_id)
+ ",
+ )
+ .try_map(|r: PgRow| {
+ assert_eq!(r.try_get_amount("in_amount", &CURRENCY)?, payment.amount);
+ assert_eq!(r.try_get_amount("bounce_amount", &CURRENCY)?, bounce_amount);
+ Ok(())
+ })
+ .fetch_one(&db)
+ .await
+ .unwrap();
+ }
+
+ #[tokio::test]
+ async fn in_simple() {
+ let (_, db) = setup().await;
+
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+
+ // Register
+ let incoming = gen_in_pay("test".to_owned());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Bounced]).await;
+
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Bounced]).await;
+
+ // Many
+ register_incoming(&db, &cfg, &gen_in_pay("another subject".to_owned()))
+ .await
+ .unwrap();
+ check_in(&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;
+
+ let original = gen_in_pay("test 2".to_owned());
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await;
+ // Idempotent
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(&db, &[Bounced, Bounced, Simple, Bounced]).await;
+ }
+
+ #[tokio::test]
+ async fn in_talerable() {
+ let (_, db) = setup().await;
+
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+ let key = EddsaPublicKey::rand();
+ let subject = format!("test with {key} reserve pub");
+
+ // Register
+ let incoming = gen_in_pay(subject.clone());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(key.clone())]).await;
+
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(key.clone())]).await;
+
+ // Key reuse is bounced
+ register_incoming(&db, &cfg, &gen_in_pay(subject.clone()))
+ .await
+ .unwrap();
+ register_incoming(&db, &cfg, &gen_in_pay(format!("another {subject}")))
+ .await
+ .unwrap();
+ check_in(&db, &[Reserve(key.clone()), Bounced, Bounced]).await;
+
+ // Admin balance adjust is ignored
+ register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
+ .await
+ .unwrap();
+ check_in(&db, &[Reserve(key.clone()), Bounced, Bounced, Simple]).await;
+
+ let new = EddsaPublicKey::rand();
+ let original = gen_in_pay(format!("test 2 with {new} reserve pub"));
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(
+ &db,
+ &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete],
+ )
+ .await;
+ // Idempotent
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(
+ &db,
+ &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete],
+ )
+ .await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(
+ &db,
+ &[Reserve(key.clone()), Bounced, Bounced, Simple, Reserve(new)],
+ )
+ .await;
+ }
+
+ #[tokio::test]
+ async fn in_mapping() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+ let first = EddsaPublicKey::rand();
+ let auth_pub = EddsaPublicKey::rand();
+ let auth_sig = EddsaSignature::rand();
+ let reference_number = subject_fmt_qr_bill(auth_pub.as_slice());
+ let subject = format!("test with MAP:{auth_pub} auth pub");
+
+ assert_eq!(
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &first,
+ &auth_pub,
+ &auth_sig,
+ false,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap(),
+ RegistrationResult::Success
+ );
+
+ // Register
+ let incoming = gen_in_pay(subject.clone());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
+
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
+
+ // Admin balance adjust is ignored
+ register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
+ .await
+ .unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple]).await;
+
+ let original = gen_in_pay(format!("test 2 for {subject}"));
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
+ // Idempotent
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
+
+ let second = EddsaPublicKey::rand();
+ assert_eq!(
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &second,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap(),
+ RegistrationResult::Success
+ );
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
+
+ // Key reuse is pending
+ for _ in 0..3 {
+ register_incoming(&db, &cfg, &gen_in_pay(subject.clone()))
+ .await
+ .unwrap();
+ }
+ check_in(
+ &db,
+ &[
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Pending,
+ Pending,
+ ],
+ )
+ .await;
+
+ // Finish pending
+ let third = EddsaPublicKey::rand();
+ assert_eq!(
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &third,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap(),
+ RegistrationResult::Success
+ );
+ check_in(
+ &db,
+ &[
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Reserve(third.clone()),
+ Pending,
+ ],
+ )
+ .await;
+ }
+
+ #[tokio::test]
+ async fn in_reference() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+ let first = EddsaPublicKey::rand();
+ let auth_pub = EddsaPublicKey::rand();
+ let auth_sig = EddsaSignature::rand();
+ let reference_number = subject_fmt_qr_bill(auth_pub.as_slice());
+
+ assert_eq!(
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &first,
+ &auth_pub,
+ &auth_sig,
+ false,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap(),
+ RegistrationResult::Success
+ );
+
+ // Register
+ let incoming = gen_in_pay(reference_number.clone());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
+
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
+
+ // Admin balance adjust is ignored
+ register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
+ .await
+ .unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple]).await;
+
+ let original = gen_in_pay(reference_number.clone());
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
+ // Idempotent
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
+
+ let second = EddsaPublicKey::rand();
+ assert_eq!(
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &second,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap(),
+ RegistrationResult::Success
+ );
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
+
+ // Key reuse is pending
+ for _ in 0..3 {
+ register_incoming(&db, &cfg, &gen_in_pay(reference_number.clone()))
+ .await
+ .unwrap();
+ }
+ check_in(
+ &db,
+ &[
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Pending,
+ Pending,
+ ],
+ )
+ .await;
+
+ // Finish pending
+ let third = EddsaPublicKey::rand();
+ assert_eq!(
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &third,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap(),
+ RegistrationResult::Success
+ );
+ check_in(
+ &db,
+ &[
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Reserve(third.clone()),
+ Pending,
+ ],
+ )
+ .await;
+ }
+
+ #[tokio::test]
+ async fn in_recover_info() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+
+ async fn check_content(db: &PgPool, p: &IncomingPayment) {
+ sqlx::query(
+ "
+ SELECT
+ uetr IS NOT DISTINCT FROM $1 AND
+ tx_id IS NOT DISTINCT FROM $2 AND
+ acct_svcr_ref IS NOT DISTINCT FROM $3 AND
+ subject IS NOT DISTINCT FROM $4 AND
+ debit_payto IS NOT DISTINCT FROM $5
+ FROM incoming_transactions ORDER BY incoming_transaction_id DESC LIMIT 1
+ ",
+ )
+ .bind(p.id.uetr)
+ .bind(&p.id.tx_id)
+ .bind(&p.id.acct_svcr_ref)
+ .bind(&p.subject)
+ .bind(p.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .try_map(|r: PgRow| {
+ assert!(r.try_get_flag(0)?);
+ Ok(())
+ })
+ .fetch_one(db)
+ .await
+ .unwrap();
+ }
+
+ // Non talerable
+ for (i, id) in [
+ IncomingId::new(Some(Uuid::new_v4()), None, None),
+ IncomingId::new(None, Some(rand_ebics_id()), None),
+ IncomingId::new(None, None, Some(rand_ebics_id())),
+ ]
+ .iter()
+ .enumerate()
+ {
+ let payment = gen_in_pay("subject".to_owned());
+
+ // Register minimal
+ let partial = IncomingPayment {
+ id: id.clone(),
+ subject: None,
+ debtor: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &partial).await.unwrap();
+ check_content(&db, &partial).await;
+ check_in_count(&db, i + 1, i, 0).await;
+
+ // Recover ID
+ let full_id = IncomingId::new(
+ Some(id.uetr.unwrap_or_else(Uuid::new_v4)),
+ Some(id.tx_id.clone().unwrap_or_else(rand_ebics_id)),
+ Some(id.acct_svcr_ref.clone().unwrap_or_else(rand_ebics_id)),
+ );
+ let full = IncomingPayment {
+ id: full_id.clone(),
+ ..partial.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 1, i, 0).await;
+
+ // Recover subject & debtor
+ let full = IncomingPayment {
+ id: full_id,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 1, i + 1, 0).await;
+ }
+
+ // Talerable
+ for (i, id) in [
+ IncomingId::new(Some(Uuid::new_v4()), None, None),
+ IncomingId::new(None, Some(rand_ebics_id()), None),
+ IncomingId::new(None, None, Some(rand_ebics_id())),
+ ]
+ .iter()
+ .enumerate()
+ {
+ let key = EddsaPublicKey::rand();
+ let payment = gen_in_pay(format!("test with {key} reserve pub"));
+
+ // Register minimal
+ let partial = IncomingPayment {
+ id: id.clone(),
+ subject: None,
+ debtor: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &partial).await.unwrap();
+ check_content(&db, &partial).await;
+ check_in_count(&db, i + 4, 3, i).await;
+
+ // Recover ID
+ let full_id = IncomingId::new(
+ Some(id.uetr.unwrap_or_else(Uuid::new_v4)),
+ Some(id.tx_id.clone().unwrap_or_else(rand_ebics_id)),
+ Some(id.acct_svcr_ref.clone().unwrap_or_else(rand_ebics_id)),
+ );
+ let full = IncomingPayment {
+ id: full_id.clone(),
+ ..partial.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 4, 3, i).await;
+
+ // Recover subject & debtor
+ let full = IncomingPayment {
+ id: full_id,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 4, 3, i + 1).await;
+ }
+ }
+
+ #[tokio::test]
+ pub async fn in_horror() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+
+ // Check we do not bounce already registered talerable transaction
+ let key = EddsaPublicKey::rand();
+ let payment = gen_in_pay(format!("test with {key} reserve pub"));
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ assert_eq!(
+ register_in_malformed(
+ &db,
+ &payment,
+ &amount("KUDOS:2.53"),
+ &rand_ebics_id(),
+ &Timestamp::now(),
+ "manual bounce",
+ )
+ .await
+ .unwrap(),
+ IncomingBounceRegistrationResult::Talerable
+ );
+ let incomplete = IncomingPayment {
+ subject: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(key.clone())]).await;
+
+ // Check we do not register as talerable bounced transaction
+ let new_key = EddsaPublicKey::rand();
+ let payment = gen_in_pay(format!("bounced {new_key}"));
+ let incomplete = IncomingPayment {
+ subject: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ check_in(&db, &[Reserve(key.clone()), Bounced]).await;
+ }
+}
diff --git a/src/lib.rs b/src/lib.rs
@@ -21,32 +21,21 @@ use std::{fmt::Display, path::Path};
use anyhow::bail;
use compact_str::CompactString;
-use jiff::Timestamp;
use rand::prelude::IndexedRandom;
use reqwest::{
Client, StatusCode,
header::{CONTENT_TYPE, HeaderValue},
};
-use sqlx::PgPool;
-use taler_api::subject::{
- IncomingSubject, parse_incoming_unstructured, parse_outgoing, subject_is_qr_bill,
-};
use taler_build::long_version;
-use taler_common::{CommonArgs, config::parser::ConfigSource, types::amount::Currency};
-use tracing::{debug, info, warn};
+use taler_common::{CommonArgs, config::parser::ConfigSource};
+use tracing::{debug, info};
use crate::{
common::EbicsLogger,
- config::{AccountType, EbicsHostCfg, NexusCfg, NexusIngestConfig},
- db::{
- InResult, IncomingBounceRegistrationResult, IncomingRegistrationResult,
- OutgoingRegistrationResult, register_in, register_in_malformed, register_in_qr_bill,
- register_in_talerable,
- },
+ config::{EbicsHostCfg, NexusCfg},
ebics_code::EbicsReturnCode,
key_management::{Order, hpb, submit_client_keys},
keys::{load_bank_keys, load_client_keys, persist_client_keys},
- model::{IncomingPayment, OutgoingBatch, OutgoingPayment},
};
use crate::{keys::ClientPriKeysFile, xml::XmlReader};
@@ -59,6 +48,7 @@ pub mod key_management;
pub mod keys;
pub mod model;
pub mod testbench;
+pub mod worker;
pub mod xml;
pub mod xml_sign;
@@ -81,240 +71,6 @@ pub struct Args {
pub common: CommonArgs,
}
-pub async fn register_incoming(
- db: &PgPool,
- cfg: &NexusIngestConfig,
- payment: &IncomingPayment,
-) -> sqlx::Result<()> {
- let log_res = |res: InResult, kind: &str, suffix: &str| {
- let fmt = std::fmt::from_fn(|f| {
- write!(f, "{payment}")?;
- if kind.is_empty() {
- write!(f, " {kind}")?;
- }
- if res.new {
- if let Some(id) = &res.bounce_id {
- write!(f, " bounced in {id}")?;
- }
- } else {
- if res.completed {
- f.write_str(" completed")?;
- if let Some(id) = &res.bounce_id {
- write!(f, " bounced in {id}")?;
- }
- } else {
- if let Some(id) = &res.bounce_id {
- write!(f, " already bounced in {id}")?;
- }
- }
- }
- if suffix.is_empty() {
- write!(f, " {suffix}")?;
- }
- Ok(())
- });
-
- if res.completed || res.new {
- info!("{fmt}")
- } else {
- debug!("{fmt}")
- }
- };
- let bounce = async |cause: &str| {
- match cfg.account_type {
- AccountType::Exchange => {
- if payment.execution_time < cfg.ignore_bounces_before {
- let res = register_in(db, payment).await?;
- log_res(res, "", &format!("ignored bounce: {cause}"));
- } else {
- let mut bounce_amount = payment.amount.clone();
- if let Some(credit_fee) = &payment.credit_fee
- && cfg.bounce_deduce_fee
- {
- if let Some(res) = bounce_amount.try_sub(credit_fee) {
- bounce_amount = res
- } else {
- let res = register_in(db, payment).await?;
- log_res(
- res,
- "",
- &format!("skip bounce (transfer fee higher than amount): {cause}"),
- );
- return Ok(());
- }
- }
- if let Some(res) = bounce_amount.try_sub(&cfg.bounce_fee) {
- bounce_amount = res
- } else {
- let res = register_in(db, payment).await?;
- log_res(
- res,
- "",
- &format!("skip bounce (bounce fee higher than amount): {cause}"),
- );
- return Ok(());
- }
- let res = register_in_malformed(
- db,
- payment,
- &bounce_amount,
- &rand_ebics_id(),
- &Timestamp::now(),
- cause,
- )
- .await?;
- match res {
- IncomingBounceRegistrationResult::Talerable => {
- warn!("{payment} tried to bounce a talerable transaction");
- }
- IncomingBounceRegistrationResult::Success(res) => {
- log_res(res, "", &format!(": {cause}"));
- }
- }
- }
- }
- AccountType::Normal => {
- let res = register_in(db, payment).await?;
- log_res(res, "", "");
- }
- }
- sqlx::Result::<_, sqlx::Error>::Ok(())
- };
-
- // Check we have enough info to handle this transaction
- if payment.debtor.is_none() {
- // TODO payment.debtor.receiverName == null
- let res = register_in(db, payment).await?;
- log_res(res, "incomplete", "");
- return Ok(());
- }
- // TODO if payment.debtor.is_none() && payment.debtor.map(|it| it.rec)
- if let Some(regex) = &cfg.restriction_payto_regex
- && let Some(debtor) = &payment.debtor
- && !regex.is_match(debtor.as_ref().as_str())
- {
- bounce("restricted account").await?;
- return Ok(());
- }
-
- if let Some(subject) = &payment.subject
- && subject_is_qr_bill(subject)
- {
- match register_in_qr_bill(db, payment, subject).await? {
- IncomingRegistrationResult::ReservePubReuse => bounce("reverse pub reuse").await?,
- IncomingRegistrationResult::MappingReuse => bounce("mapping reuse").await?,
- IncomingRegistrationResult::UnknownMapping => bounce("unknown mapping").await?,
- IncomingRegistrationResult::Success(res) => {
- log_res(res, "", "");
- }
- }
- } else {
- match parse_incoming_unstructured(payment.subject.as_deref().unwrap_or_default()) {
- Ok(None) => bounce("missing public key").await?,
- Ok(Some(IncomingSubject::AdminBalanceAdjust)) => {
- let res = register_in(db, payment).await?;
- log_res(res, "admin balance adjust", "");
- }
- Ok(Some(subject)) => match register_in_talerable(db, payment, &subject).await? {
- IncomingRegistrationResult::ReservePubReuse => bounce("reverse pub reuse").await?,
- IncomingRegistrationResult::MappingReuse => bounce("mapping reuse").await?,
- IncomingRegistrationResult::UnknownMapping => bounce("unknown mapping").await?,
- IncomingRegistrationResult::Success(res) => {
- log_res(res, "", "");
- }
- },
- Err(e) => {
- bounce(&e.to_string()).await?;
- }
- }
- }
-
- Ok(())
-}
-
-pub async fn register_outgoing(
- db: &PgPool,
- payment: &OutgoingPayment,
-) -> sqlx::Result<OutgoingRegistrationResult> {
- let metadata = payment
- .subject
- .as_ref()
- .and_then(|s| parse_outgoing(s).ok());
- let res = db::register_out_tx(db, payment, metadata.as_ref()).await?;
- if res.new {
- if res.initiated {
- info!("{payment}");
- } else {
- warn!("{payment} recovered");
- }
- } else {
- debug!("{payment} already seen");
- }
- Ok(res)
-}
-
-pub async fn register_outgoing_batch(
- db: &PgPool,
- currency: &Currency,
- batch: &OutgoingBatch,
-) -> sqlx::Result<()> {
- info!("{batch}");
- let txs = db::unsettled_tx_in_batch(db, currency, &batch.msg_id, &batch.execution_time).await?;
- for tx in txs {
- register_outgoing(db, &tx).await?;
- }
- Ok(())
-}
-
-pub enum Tx {
- In(IncomingPayment),
- Out(OutgoingPayment),
- Batch(OutgoingBatch),
- Reversal,
-}
-
-impl Tx {
- pub fn execution_time(&self) -> &Timestamp {
- match self {
- Tx::In(IncomingPayment { execution_time, .. })
- | Tx::Out(OutgoingPayment { execution_time, .. })
- | Tx::Batch(OutgoingBatch { execution_time, .. }) => execution_time,
- Tx::Reversal => todo!(),
- }
- }
-}
-
-impl Display for Tx {
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
- match self {
- Tx::In(incoming_payment) => incoming_payment.fmt(f),
- Tx::Out(outgoing_payment) => outgoing_payment.fmt(f),
- Tx::Batch(outgoing_batch) => outgoing_batch.fmt(f),
- Tx::Reversal => todo!(),
- }
- }
-}
-
-pub async fn register_tx(db: &PgPool, cfg: &NexusIngestConfig, tx: &Tx) -> sqlx::Result<()> {
- if tx.execution_time() < &cfg.ignore_txs_before {
- debug!("IGNORE {tx}");
- } else {
- match tx {
- Tx::In(payment) => {
- register_incoming(db, cfg, payment).await?;
- }
- Tx::Out(payment) => {
- register_outgoing(db, payment).await?;
- }
- Tx::Batch(batch) => {
- register_outgoing_batch(db, &cfg.currency, batch).await?;
- }
- Tx::Reversal => todo!(),
- }
- }
- Ok(())
-}
-
/** Load client private keys at or create new ones if missing */
pub fn load_or_generate_client_keys(path: &Path) -> anyhow::Result<ClientPriKeysFile> {
// If exists load from disk
diff --git a/src/model.rs b/src/model.rs
@@ -347,3 +347,32 @@ pub struct InitiatedPayment {
pub initiation_time: Timestamp,
pub end_to_end_id: CompactString,
}
+
+pub enum Tx {
+ In(IncomingPayment),
+ Out(OutgoingPayment),
+ Batch(OutgoingBatch),
+ Reversal,
+}
+
+impl Tx {
+ pub fn execution_time(&self) -> &Timestamp {
+ match self {
+ Tx::In(IncomingPayment { execution_time, .. })
+ | Tx::Out(OutgoingPayment { execution_time, .. })
+ | Tx::Batch(OutgoingBatch { execution_time, .. }) => execution_time,
+ Tx::Reversal => todo!(),
+ }
+ }
+}
+
+impl Display for Tx {
+ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ match self {
+ Tx::In(incoming_payment) => incoming_payment.fmt(f),
+ Tx::Out(outgoing_payment) => outgoing_payment.fmt(f),
+ Tx::Batch(outgoing_batch) => outgoing_batch.fmt(f),
+ Tx::Reversal => todo!(),
+ }
+ }
+}
diff --git a/src/worker.rs b/src/worker.rs
@@ -0,0 +1,245 @@
+/*
+* 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 sqlx::PgPool;
+use taler_api::subject::{
+ IncomingSubject, parse_incoming_unstructured, parse_outgoing, subject_is_qr_bill,
+};
+use taler_common::types::amount::Currency;
+use tracing::{debug, info, warn};
+
+use crate::{
+ config::{AccountType, NexusIngestConfig},
+ db::{
+ initiated::unsettled_tx_in_batch,
+ payment::{
+ InResult, IncomingBounceRegistrationResult, IncomingRegistrationResult,
+ OutgoingRegistrationResult, register_in, register_in_malformed, register_in_qr_bill,
+ register_in_talerable, register_out_tx,
+ },
+ },
+ model::{IncomingPayment, OutgoingBatch, OutgoingPayment, Tx},
+ rand_ebics_id,
+};
+
+pub async fn register_incoming(
+ db: &PgPool,
+ cfg: &NexusIngestConfig,
+ payment: &IncomingPayment,
+) -> sqlx::Result<()> {
+ let log_res = |res: InResult, kind: &str, suffix: &str| {
+ let fmt = std::fmt::from_fn(|f| {
+ write!(f, "{payment}")?;
+ if kind.is_empty() {
+ write!(f, " {kind}")?;
+ }
+ if res.new {
+ if let Some(id) = &res.bounce_id {
+ write!(f, " bounced in {id}")?;
+ }
+ } else {
+ if res.completed {
+ f.write_str(" completed")?;
+ if let Some(id) = &res.bounce_id {
+ write!(f, " bounced in {id}")?;
+ }
+ } else {
+ if let Some(id) = &res.bounce_id {
+ write!(f, " already bounced in {id}")?;
+ }
+ }
+ }
+ if suffix.is_empty() {
+ write!(f, " {suffix}")?;
+ }
+ Ok(())
+ });
+
+ if res.completed || res.new {
+ info!("{fmt}")
+ } else {
+ debug!("{fmt}")
+ }
+ };
+ let bounce = async |cause: &str| {
+ match cfg.account_type {
+ AccountType::Exchange => {
+ if payment.execution_time < cfg.ignore_bounces_before {
+ let res = register_in(db, payment).await?;
+ log_res(res, "", &format!("ignored bounce: {cause}"));
+ } else {
+ let mut bounce_amount = payment.amount.clone();
+ if let Some(credit_fee) = &payment.credit_fee
+ && cfg.bounce_deduce_fee
+ {
+ if let Some(res) = bounce_amount.try_sub(credit_fee) {
+ bounce_amount = res
+ } else {
+ let res = register_in(db, payment).await?;
+ log_res(
+ res,
+ "",
+ &format!("skip bounce (transfer fee higher than amount): {cause}"),
+ );
+ return Ok(());
+ }
+ }
+ if let Some(res) = bounce_amount.try_sub(&cfg.bounce_fee) {
+ bounce_amount = res
+ } else {
+ let res = register_in(db, payment).await?;
+ log_res(
+ res,
+ "",
+ &format!("skip bounce (bounce fee higher than amount): {cause}"),
+ );
+ return Ok(());
+ }
+ let res = register_in_malformed(
+ db,
+ payment,
+ &bounce_amount,
+ &rand_ebics_id(),
+ &Timestamp::now(),
+ cause,
+ )
+ .await?;
+ match res {
+ IncomingBounceRegistrationResult::Talerable => {
+ warn!("{payment} tried to bounce a talerable transaction");
+ }
+ IncomingBounceRegistrationResult::Success(res) => {
+ log_res(res, "", &format!(": {cause}"));
+ }
+ }
+ }
+ }
+ AccountType::Normal => {
+ let res = register_in(db, payment).await?;
+ log_res(res, "", "");
+ }
+ }
+ sqlx::Result::<_, sqlx::Error>::Ok(())
+ };
+
+ // Check we have enough info to handle this transaction
+ if payment.debtor.is_none() {
+ // TODO payment.debtor.receiverName == null
+ let res = register_in(db, payment).await?;
+ log_res(res, "incomplete", "");
+ return Ok(());
+ }
+ // TODO if payment.debtor.is_none() && payment.debtor.map(|it| it.rec)
+ if let Some(regex) = &cfg.restriction_payto_regex
+ && let Some(debtor) = &payment.debtor
+ && !regex.is_match(debtor.as_ref().as_str())
+ {
+ bounce("restricted account").await?;
+ return Ok(());
+ }
+
+ if let Some(subject) = &payment.subject
+ && subject_is_qr_bill(subject)
+ {
+ match register_in_qr_bill(db, payment, subject).await? {
+ IncomingRegistrationResult::ReservePubReuse => bounce("reverse pub reuse").await?,
+ IncomingRegistrationResult::MappingReuse => bounce("mapping reuse").await?,
+ IncomingRegistrationResult::UnknownMapping => bounce("unknown mapping").await?,
+ IncomingRegistrationResult::Success(res) => {
+ log_res(res, "", "");
+ }
+ }
+ } else {
+ match parse_incoming_unstructured(payment.subject.as_deref().unwrap_or_default()) {
+ Ok(None) => bounce("missing public key").await?,
+ Ok(Some(IncomingSubject::AdminBalanceAdjust)) => {
+ let res = register_in(db, payment).await?;
+ log_res(res, "admin balance adjust", "");
+ }
+ Ok(Some(subject)) => match register_in_talerable(db, payment, &subject).await? {
+ IncomingRegistrationResult::ReservePubReuse => bounce("reverse pub reuse").await?,
+ IncomingRegistrationResult::MappingReuse => bounce("mapping reuse").await?,
+ IncomingRegistrationResult::UnknownMapping => bounce("unknown mapping").await?,
+ IncomingRegistrationResult::Success(res) => {
+ log_res(res, "", "");
+ }
+ },
+ Err(e) => {
+ bounce(&e.to_string()).await?;
+ }
+ }
+ }
+
+ Ok(())
+}
+
+pub async fn register_outgoing(
+ db: &PgPool,
+ payment: &OutgoingPayment,
+) -> sqlx::Result<OutgoingRegistrationResult> {
+ let metadata = payment
+ .subject
+ .as_ref()
+ .and_then(|s| parse_outgoing(s).ok());
+ let res = register_out_tx(db, payment, metadata.as_ref()).await?;
+ if res.new {
+ if res.initiated {
+ info!("{payment}");
+ } else {
+ warn!("{payment} recovered");
+ }
+ } else {
+ debug!("{payment} already seen");
+ }
+ Ok(res)
+}
+
+pub async fn register_outgoing_batch(
+ db: &PgPool,
+ currency: &Currency,
+ batch: &OutgoingBatch,
+) -> sqlx::Result<()> {
+ info!("{batch}");
+ let txs = unsettled_tx_in_batch(db, currency, &batch.msg_id, &batch.execution_time).await?;
+ for tx in txs {
+ register_outgoing(db, &tx).await?;
+ }
+ Ok(())
+}
+
+pub async fn register_tx(db: &PgPool, cfg: &NexusIngestConfig, tx: &Tx) -> sqlx::Result<()> {
+ if tx.execution_time() < &cfg.ignore_txs_before {
+ debug!("IGNORE {tx}");
+ } else {
+ match tx {
+ Tx::In(payment) => {
+ register_incoming(db, cfg, payment).await?;
+ }
+ Tx::Out(payment) => {
+ register_outgoing(db, payment).await?;
+ }
+ Tx::Batch(batch) => {
+ register_outgoing_batch(db, &cfg.currency, batch).await?;
+ }
+ Tx::Reversal => todo!(),
+ }
+ }
+ Ok(())
+}