libeufin

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

initiated.rs (32139B)


      1 /*
      2   This file is part of TALER
      3   Copyright (C) 2026 Taler Systems SA
      4 
      5   TALER is free software; you can redistribute it and/or modify it under the
      6   terms of the GNU Affero General Public License as published by the Free Software
      7   Foundation; either version 3, or (at your option) any later version.
      8 
      9   TALER is distributed in the hope that it will be useful, but WITHOUT ANY
     10   WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR
     11   A PARTICULAR PURPOSE.  See the GNU Affero General Public License for more details.
     12 
     13   You should have received a copy of the GNU Affero General Public License along with
     14   TALER; see the file COPYING.  If not, see <http://www.gnu.org/licenses/>
     15 */
     16 
     17 use std::collections::BTreeMap;
     18 
     19 use const_format::formatcp;
     20 use jiff::Timestamp;
     21 use libeufin_ebics::iso20022::model::{OutId, OutTx};
     22 use sqlx::{PgPool, Row as _, postgres::PgRow};
     23 use taler_api::{
     24     db::{BindHelper as _, TypeHelper as _},
     25     serialized,
     26 };
     27 use taler_common::types::{
     28     amount::{Amount, Currency},
     29     payto::FullIbanPayto,
     30 };
     31 
     32 use crate::{
     33     db::{PENDING, UNSETTLED},
     34     model::{Initiated, PaymentBatch, SubmissionState},
     35 };
     36 
     37 /// Outgoing payments initiation result
     38 #[derive(Debug, PartialEq, Eq)]
     39 pub enum PaymentInitiationResult {
     40     Success(u64),
     41     RequestUidReuse,
     42 }
     43 
     44 /// Initiate a new payment
     45 pub async fn initiate(
     46     pool: &PgPool,
     47     amount: &Amount,
     48     subject: &str,
     49     creditor: &FullIbanPayto,
     50     initiation_time: &Timestamp,
     51     e2e_id: &str,
     52 ) -> sqlx::Result<PaymentInitiationResult> {
     53     let res = serialized!(
     54         sqlx::query(
     55             "
     56         INSERT INTO initiated_outgoing_transactions (
     57             amount,
     58             subject,
     59             credit_payto,
     60             initiation_time,
     61             end_to_end_id
     62         ) VALUES ($1,$2,$3,$4,$5)
     63         RETURNING initiated_outgoing_transaction_id
     64         ",
     65         )
     66         .bind(amount)
     67         .bind(subject)
     68         .bind(creditor.as_uri().as_ref().as_str())
     69         .bind_timestamp(initiation_time)
     70         .bind(e2e_id)
     71         .try_map(|r: PgRow| Ok(PaymentInitiationResult::Success(r.try_get_u64(0)?)))
     72         .fetch_one(pool)
     73     );
     74     if let Err(e) = &res
     75         && let Some(db_err) = e.as_database_error()
     76         && db_err.code() == Some(std::borrow::Cow::Borrowed("23505"))
     77     {
     78         Ok(PaymentInitiationResult::RequestUidReuse)
     79     } else {
     80         res
     81     }
     82 }
     83 
     84 /// Group unbatched transaction into a single batch
     85 pub async fn batch_initiated(
     86     pool: &PgPool,
     87     timestamp: &Timestamp,
     88     ebics_id: &str,
     89     require_ack: bool,
     90 ) -> sqlx::Result<()> {
     91     serialized!(
     92         sqlx::query("SELECT batch_outgoing_transactions($1, $2, $3)")
     93             .bind_timestamp(timestamp)
     94             .bind(ebics_id)
     95             .bind(require_ack)
     96             .execute(pool)
     97     )?;
     98     Ok(())
     99 }
    100 
    101 pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<bool> {
    102     let res = serialized!(
    103         sqlx::query("UPDATE initiated_outgoing_transactions SET awaiting_ack=false WHERE initiated_outgoing_transaction_id=$1")
    104     .bind(id as i64)
    105     .execute(db)
    106     )?;
    107     Ok(res.rows_affected() > 0)
    108 }
    109 
    110 pub async fn initiated_submittable(
    111     db: &PgPool,
    112     currency: &Currency,
    113 ) -> sqlx::Result<Vec<PaymentBatch>> {
    114     const SELECT_PART: &str = "
    115         SELECT initiated_outgoing_batch_id, message_id, creation_date, sum
    116         FROM initiated_outgoing_batches
    117     ";
    118     serialized!(async {
    119         let mut tx = db.begin().await?;
    120         // We want to maximize the number of successfully submitted batches in the event
    121         // of a malformed transaction or a persistent error classified as transient. We send
    122         // the unsubmitted batches first, starting with the oldest by creation time.
    123         // This is the happy path, giving every batch a chance while being fair on the
    124         // basis of creation date.
    125         // Then we retry the failed batches, starting with the oldest by submission time.
    126         // This the bad path retrying each failed batch applying a rotation based on
    127         // resubmission time.
    128         let mut batches = sqlx::query(formatcp!(
    129             "
    130         ({SELECT_PART} WHERE status='unsubmitted' ORDER BY creation_date ASC)
    131             UNION ALL
    132         ({SELECT_PART} WHERE status='transient_failure' ORDER BY submission_date)
    133             "
    134         ))
    135         .try_map(|r: PgRow| {
    136             Ok(PaymentBatch {
    137                 id: r.try_get_u64("initiated_outgoing_batch_id")?,
    138                 msg_id: r.try_get("message_id")?,
    139                 creation_date: r.try_get_timestamp("creation_date")?,
    140                 sum: r.try_get_amount("sum", currency)?,
    141                 payments: Vec::new(),
    142             })
    143         })
    144         .fetch_all(&mut *tx)
    145         .await?;
    146         let mut batch_map: BTreeMap<_, _> = batches.iter_mut().map(|it| (it.id, it)).collect();
    147         // Then load transactions
    148         sqlx::query(
    149             "
    150         SELECT
    151             initiated_outgoing_transaction_id
    152             ,amount
    153             ,subject
    154             ,credit_payto
    155             ,initiated_outgoing_transactions.initiation_time
    156             ,end_to_end_id
    157             ,initiated_outgoing_batch_id
    158         FROM initiated_outgoing_transactions
    159             JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
    160         WHERE initiated_outgoing_batches.status IN ('unsubmitted', 'transient_failure')
    161         ",
    162         )
    163         .try_map(|r: PgRow| {
    164             let payment = Initiated {
    165                 id: r.try_get_u64("initiated_outgoing_transaction_id")?,
    166                 amount: r.try_get_amount("amount", currency)?,
    167                 creditor: r.try_get_parse("credit_payto")?,
    168                 subject: r.try_get("subject")?,
    169                 initiation_time: r.try_get_timestamp("initiation_time")?,
    170                 e2e_id: r.try_get("end_to_end_id")?,
    171             };
    172             let batch_id = r.try_get_u64("initiated_outgoing_batch_id")?;
    173             batch_map.get_mut(&batch_id).unwrap().payments.push(payment);
    174             Ok(())
    175         })
    176         .fetch_all(&mut *tx)
    177         .await?;
    178         tx.commit().await?;
    179         Ok(batches)
    180     })
    181 }
    182 
    183 pub async fn unsettled_tx_in_batch(
    184     db: &PgPool,
    185     currency: &Currency,
    186     msg_id: &str,
    187     execution_time: &Timestamp,
    188 ) -> sqlx::Result<Vec<OutTx>> {
    189     serialized!(
    190         sqlx::query(formatcp!(
    191             "
    192         SELECT
    193             end_to_end_id,
    194             amount,
    195             subject,
    196             credit_payto
    197         FROM initiated_outgoing_transactions
    198             JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
    199         WHERE message_id = $1
    200             AND initiated_outgoing_transactions.{UNSETTLED}
    201         "
    202         ))
    203         .bind(msg_id)
    204         .try_map(|r: PgRow| {
    205             Ok(OutTx {
    206                 id: OutId {
    207                     msg_id: Some(msg_id.into()),
    208                     e2e_id: r.try_get("end_to_end_id")?,
    209                     sref: None,
    210                 },
    211                 amount: r.try_get_amount("amount", currency)?,
    212                 debit_fee: Amount::zero(currency),
    213                 subject: r.try_get("subject")?,
    214                 execution_time: *execution_time,
    215                 creditor: r.try_get_opt_payto("credit_payto")?,
    216             })
    217         })
    218         .fetch_all(db)
    219     )
    220 }
    221 
    222 /** Register submission success of order [orderId] for batch [id] at [timestamp] */
    223 pub async fn batch_sub_success(
    224     db: &PgPool,
    225     batch_id: u64,
    226     timestamp: &Timestamp,
    227     order_id: &str,
    228 ) -> sqlx::Result<()> {
    229     serialized!(async {
    230         let mut tx = db.begin().await?;
    231         // Update batch status
    232         let updated = sqlx::query(
    233             "
    234             UPDATE initiated_outgoing_batches
    235             SET  status = 'pending'
    236                 ,submission_date = $1
    237                 ,status_msg = NULL
    238                 ,order_id = $2
    239                 ,submission_counter = submission_counter + 1
    240             WHERE initiated_outgoing_batch_id = $3 AND order_id IS NULL
    241             ",
    242         )
    243         .bind_timestamp(timestamp)
    244         .bind(order_id)
    245         .bind(batch_id as i64)
    246         .execute(&mut *tx)
    247         .await?;
    248         if updated.rows_affected() > 0 {
    249             // Update unsettled batch's transaction status
    250             sqlx::query(formatcp!(
    251                 "
    252             UPDATE initiated_outgoing_transactions
    253             SET status = 'pending', status_msg = NULL
    254             WHERE initiated_outgoing_batch_id = $1 AND {UNSETTLED}
    255             "
    256             ))
    257             .bind(batch_id as i64)
    258             .execute(&mut *tx)
    259             .await?;
    260         }
    261         tx.commit().await
    262     })
    263 }
    264 
    265 /** Register submission failure with [msg] for batch [id] at [timestamp]*/
    266 pub async fn batch_sub_failure(
    267     db: &PgPool,
    268     batch_id: u64,
    269     timestamp: &Timestamp,
    270     msg: &str,
    271 ) -> sqlx::Result<()> {
    272     let permanent = false;
    273     serialized!(async {
    274         let mut tx = db.begin().await?;
    275         // Update batch status
    276         sqlx::query(
    277             "
    278             UPDATE initiated_outgoing_batches
    279             SET  status = $1
    280                 ,submission_date = $2
    281                 ,status_msg = $3
    282                 ,submission_counter = submission_counter + 1
    283             WHERE initiated_outgoing_batch_id = $4
    284             ",
    285         )
    286         .bind(if permanent {
    287             SubmissionState::permanent_failure
    288         } else {
    289             SubmissionState::transient_failure
    290         })
    291         .bind_timestamp(timestamp)
    292         .bind(msg)
    293         .bind(batch_id as i64)
    294         .execute(&mut *tx)
    295         .await?;
    296         // Update unsettled batch's transaction status
    297         sqlx::query(formatcp!(
    298             "
    299             UPDATE initiated_outgoing_transactions
    300             SET status = $1, status_msg = $2
    301             WHERE initiated_outgoing_batch_id = $3 AND {UNSETTLED}
    302             "
    303         ))
    304         .bind(if permanent {
    305             SubmissionState::permanent_failure
    306         } else {
    307             SubmissionState::transient_failure
    308         })
    309         .bind(msg)
    310         .bind(batch_id as i64)
    311         .execute(&mut *tx)
    312         .await?;
    313         tx.commit().await
    314     })
    315 }
    316 
    317 /** Register order step [msg] for [orderId] */
    318 pub async fn order_step(db: &PgPool, order_id: &str, msg: &str) -> sqlx::Result<()> {
    319     serialized!(async {
    320         let mut tx = db.begin().await?;
    321         // Update batch status
    322         let batch_id = sqlx::query(formatcp!(
    323             "
    324             UPDATE initiated_outgoing_batches
    325             SET status = 'pending', status_msg = $1
    326             WHERE order_id = $2 AND {PENDING}
    327             RETURNING initiated_outgoing_batch_id
    328             "
    329         ))
    330         .bind(msg)
    331         .bind(order_id)
    332         .try_map(|r: PgRow| r.try_get_u64(0))
    333         .fetch_optional(&mut *tx)
    334         .await?;
    335         if let Some(batch_id) = batch_id {
    336             // Update unsettled batch's transaction status
    337             sqlx::query(formatcp!(
    338                 "
    339             UPDATE initiated_outgoing_transactions
    340             SET status = 'pending', status_msg = $1
    341             WHERE initiated_outgoing_batch_id = $2 AND {PENDING}
    342             "
    343             ))
    344             .bind(msg)
    345             .bind(batch_id as i64)
    346             .execute(&mut *tx)
    347             .await?;
    348         }
    349         tx.commit().await
    350     })
    351 }
    352 
    353 /** Register order success for [orderId] and return message_id if found */
    354 pub async fn order_success(db: &PgPool, order_id: &str) -> sqlx::Result<Option<String>> {
    355     serialized!(async {
    356         let mut tx = db.begin().await?;
    357         // Update batch status
    358         let res = sqlx::query(formatcp!(
    359             "
    360         UPDATE initiated_outgoing_batches
    361         SET status = 'success'
    362         WHERE order_id = $1
    363         RETURNING initiated_outgoing_batch_id, message_id
    364         "
    365         ))
    366         .bind(order_id)
    367         .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?)))
    368         .fetch_optional(&mut *tx)
    369         .await?;
    370         if let Some((batch_id, _)) = &res {
    371             // Update unsettled batch's transaction status
    372             sqlx::query(formatcp!(
    373                 "
    374             UPDATE initiated_outgoing_transactions
    375             SET status = 'pending'
    376             WHERE initiated_outgoing_batch_id = $1 AND {UNSETTLED}
    377             "
    378             ))
    379             .bind(*batch_id as i64)
    380             .execute(&mut *tx)
    381             .await?;
    382         }
    383         tx.commit().await?;
    384         Ok(res.map(|(_, msg_id)| msg_id))
    385     })
    386 }
    387 
    388 /** Register order failure for [orderId] and return message_id and previous status_msg if found */
    389 pub async fn order_failure(
    390     db: &PgPool,
    391     order_id: &str,
    392 ) -> sqlx::Result<Option<(String, Option<String>)>> {
    393     serialized!(async {
    394         let mut tx = db.begin().await?;
    395         // Update batch status
    396         let res = sqlx::query(formatcp!(
    397             "
    398         UPDATE initiated_outgoing_batches
    399         SET status = 'permanent_failure'
    400         WHERE order_id = $1
    401         RETURNING initiated_outgoing_batch_id, message_id, status_msg
    402         "
    403         ))
    404         .bind(order_id)
    405         .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?, r.try_get(2)?)))
    406         .fetch_optional(&mut *tx)
    407         .await?;
    408         if let Some((batch_id, _, _)) = &res {
    409             // Update unsettled batch's transaction status
    410             sqlx::query(formatcp!(
    411                 "
    412             UPDATE initiated_outgoing_transactions
    413             SET status = 'permanent_failure'
    414             WHERE initiated_outgoing_batch_id = $1
    415             "
    416             ))
    417             .bind(*batch_id as i64)
    418             .execute(&mut *tx)
    419             .await?;
    420         }
    421         tx.commit().await?;
    422         Ok(res.map(|(_, msg_id, status_msg)| (msg_id, status_msg)))
    423     })
    424 }
    425 
    426 /** Register payment status [state] with [msg] for batch [msgId] */
    427 pub async fn batch_status_update(
    428     db: &PgPool,
    429     msg_id: &str,
    430     state: SubmissionState,
    431     msg: &str,
    432 ) -> sqlx::Result<bool> {
    433     serialized!(
    434         sqlx::query(formatcp!(
    435             "SELECT out_ok FROM batch_status_update($1,$2,$3)"
    436         ))
    437         .bind(msg_id)
    438         .bind(state)
    439         .bind((!msg.is_empty()).then_some(msg))
    440         .try_map(|r: PgRow| r.try_get(0))
    441         .fetch_one(db)
    442     )
    443 }
    444 
    445 /** Register payment status [state] with [msg] for transaction [endToEndId] in batch [msgId] */
    446 pub async fn tx_status_update(
    447     db: &PgPool,
    448     end_to_end_id: &str,
    449     msg_id: &str,
    450     state: SubmissionState,
    451     msg: &str,
    452 ) -> sqlx::Result<bool> {
    453     serialized!(
    454         sqlx::query(formatcp!(
    455             "SELECT out_ok FROM tx_status_update($1,$2,$3,$4)"
    456         ))
    457         .bind(end_to_end_id)
    458         .bind(msg_id)
    459         .bind(state)
    460         .bind(msg)
    461         .try_map(|r: PgRow| r.try_get(0))
    462         .fetch_one(db)
    463     )
    464 }
    465 
    466 #[cfg(test)]
    467 mod test {
    468     use std::str::FromStr as _;
    469 
    470     use jiff::{Span, Timestamp, civil::Date};
    471     use libeufin_ebics::{ebics::rand_ebics_id, iso20022::model::Tx};
    472     use sqlx::{PgPool, Row as _, postgres::PgRow};
    473     use taler_api::db::TypeHelper as _;
    474     use taler_common::{config::Config, types::utils::date_to_utc_ts};
    475     use taler_macros::db_test;
    476 
    477     use crate::{
    478         config::{NexusCfg, NexusIngestCfg},
    479         constants::CONFIG_SOURCE,
    480         db::{
    481             initiated::{
    482                 PaymentInitiationResult, batch_initiated, batch_status_update, batch_sub_failure,
    483                 batch_sub_success, initiated_submittable, order_failure, order_step, order_success,
    484                 tx_status_update,
    485             },
    486             test::check_count,
    487         },
    488         fetch::{register_outgoing, register_tx},
    489         model::SubmissionState,
    490         test::{CURR, gen_in_pay, gen_initiate, gen_out_pay},
    491     };
    492 
    493     #[db_test]
    494     pub async fn initiated_skip(db: PgPool) {
    495         let cfg = Config::load(CONFIG_SOURCE, Some("conf/skip.conf")).unwrap();
    496         let cfg = NexusCfg::parse(&cfg).unwrap();
    497         let cfg = cfg.ingest().unwrap();
    498         let millis = Span::new().milliseconds(10);
    499 
    500         async fn ingest(db: &PgPool, cfg: &NexusIngestCfg, execution_time: Timestamp) {
    501             for tx in [
    502                 Tx::In(
    503                     gen_in_pay(format_args!("test at {execution_time}"))
    504                         .with_execution_time(execution_time),
    505                 ),
    506                 Tx::Out(
    507                     gen_out_pay(format!("test at {execution_time}"))
    508                         .with_execution_time(execution_time),
    509                 ),
    510             ] {
    511                 register_tx(db, cfg, &tx).await.unwrap()
    512             }
    513         }
    514 
    515         assert_eq!(
    516             cfg.ignore_txs_before,
    517             date_to_utc_ts(&Date::from_str("2024-04-04").unwrap())
    518         );
    519         assert_eq!(
    520             cfg.ignore_bounces_before,
    521             date_to_utc_ts(&Date::from_str("2024-06-12").unwrap())
    522         );
    523 
    524         // No transaction at the beginning
    525         check_count(&db, 0, 0).await;
    526 
    527         // Skipped transactions
    528         ingest(&db, &cfg, cfg.ignore_txs_before - millis).await;
    529         check_count(&db, 0, 0).await;
    530 
    531         // Skipped bounces
    532         ingest(&db, &cfg, cfg.ignore_txs_before).await;
    533         ingest(&db, &cfg, cfg.ignore_txs_before + millis).await;
    534         ingest(&db, &cfg, cfg.ignore_bounces_before - millis).await;
    535         check_count(&db, 6, 0).await;
    536 
    537         // Bounces
    538         ingest(&db, &cfg, cfg.ignore_bounces_before).await;
    539         ingest(&db, &cfg, cfg.ignore_bounces_before + millis).await;
    540         check_count(&db, 10, 2).await;
    541     }
    542 
    543     #[db_test]
    544     pub async fn initiated_status(db: PgPool) {
    545         use SubmissionState::*;
    546 
    547         let check_parts = async |batch_id: u64,
    548                                  batch_status: SubmissionState,
    549                                  batch_msg: &str,
    550                                  tx_status: SubmissionState,
    551                                  tx_msg: &str,
    552                                  settled_status: SubmissionState,
    553                                  settled_msg: &str| {
    554             // Check batch status
    555             let msg_id: String = sqlx::query(
    556                 "
    557                 SELECT message_id, status, status_msg FROM initiated_outgoing_batches WHERE initiated_outgoing_batch_id=$1
    558                 "
    559             ).bind(batch_id as i64)
    560             .try_map(|r: PgRow| {
    561                 let msg_id: String =  r.try_get("message_id")?;
    562                 assert_eq!((batch_status, Some(batch_msg).filter(|it| !it.is_empty())), (r.try_get("status")?, r.try_get("status_msg")?), "{msg_id}");
    563                 Ok(msg_id)
    564             }).fetch_one(&db).await.unwrap();
    565             // Check tx status
    566             sqlx::query(
    567                 "
    568                 SELECT end_to_end_id, status, status_msg FROM initiated_outgoing_transactions WHERE initiated_outgoing_batch_id=$1
    569                 "
    570             ).bind(batch_id as i64).try_map(|r: PgRow| {
    571                  let end_to_end_id: &str = r.try_get("end_to_end_id")?;
    572                     let expected = match end_to_end_id {
    573                         "TX" => (tx_status, Some(tx_msg).filter(|it| !it.is_empty())),
    574                         "TX_SETTLED" => (settled_status, Some(settled_msg).filter(|it| !it.is_empty())),
    575                         _ =>panic!("Unexpected tx $endToEndId")
    576                     };
    577                     assert_eq!(expected,
    578                         (r.try_get("status")?, r.try_get("status_msg")?),
    579                         "{msg_id},{end_to_end_id}"
    580                     );
    581                 Ok(())
    582             }).fetch_all(&db).await.unwrap();
    583         };
    584 
    585         let check_batch_tx = async |batch_id: u64,
    586                                     status: SubmissionState,
    587                                     msg: &str,
    588                                     tx_status: SubmissionState| {
    589             check_parts(batch_id, status, msg, tx_status, msg, tx_status, msg).await;
    590         };
    591         let check_batch = async |batch_id: u64, status: SubmissionState, msg: &str| {
    592             check_batch_tx(batch_id, status, msg, status).await;
    593         };
    594         let check_order_tx = async |order_id: &str,
    595                                     status: SubmissionState,
    596                                     msg: &str,
    597                                     tx_status: SubmissionState| {
    598             let batch_id = sqlx::query(
    599                 "SELECT initiated_outgoing_batch_id FROM initiated_outgoing_batches WHERE order_id=$1"
    600             ).bind(order_id)
    601             .try_map(|r: PgRow| {
    602                 r.try_get_u64(0)
    603             }).fetch_one(&db).await.unwrap();
    604             check_batch_tx(batch_id, status, msg, tx_status).await;
    605         };
    606         let check_order = async |order_id: &str, status: SubmissionState, msg: &str| {
    607             check_order_tx(order_id, status, msg, status).await;
    608         };
    609 
    610         async fn test(db: &PgPool, lambda: impl AsyncFnOnce(u64)) {
    611             // Reset DB
    612             sqlx::query("DELETE FROM initiated_outgoing_transactions")
    613                 .execute(db)
    614                 .await
    615                 .unwrap();
    616             sqlx::query("DELETE FROM initiated_outgoing_batches")
    617                 .execute(db)
    618                 .await
    619                 .unwrap();
    620             // Create a test batch with three transactions
    621             for id in ["TX", "TX_SETTLED"] {
    622                 assert!(matches!(
    623                     gen_initiate(db, id, "lol").await,
    624                     PaymentInitiationResult::Success(_)
    625                 ));
    626             }
    627             batch_initiated(db, &Timestamp::now(), "BATCH", false)
    628                 .await
    629                 .unwrap();
    630 
    631             // Create witness transactions and batch
    632             for id in ["WITNESS_1", "WITNESS_2"] {
    633                 assert!(matches!(
    634                     gen_initiate(db, id, "lol").await,
    635                     PaymentInitiationResult::Success(_)
    636                 ));
    637             }
    638             batch_initiated(db, &Timestamp::now(), "BATCH_WITNESS", false)
    639                 .await
    640                 .unwrap();
    641             for id in ["WITNESS_3", "WITNESS_4"] {
    642                 assert!(matches!(
    643                     gen_initiate(db, id, "lol").await,
    644                     PaymentInitiationResult::Success(_)
    645                 ));
    646             }
    647             // Check everything is unsubmitted
    648             sqlx::query(
    649                 "
    650                 SELECT (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_batches)
    651                    AND (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_transactions)
    652                 "
    653             ).try_map(|r: PgRow| {
    654                 assert!(r.try_get_flag(0).unwrap());
    655                 Ok(())
    656             }).fetch_one(db).await.unwrap();
    657             let submitibale = initiated_submittable(db, &CURR).await.unwrap();
    658             lambda(
    659                 submitibale
    660                     .iter()
    661                     .find(|it| it.msg_id == "BATCH")
    662                     .unwrap()
    663                     .id,
    664             )
    665             .await;
    666             // Check witness status is unaltered
    667             sqlx::query(
    668                 "
    669                 SELECT (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_batches WHERE message_id != 'BATCH')
    670                    AND (SELECT bool_and(initiated_outgoing_transactions.status = 'unsubmitted')
    671                             FROM initiated_outgoing_transactions JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id)
    672                             WHERE message_id != 'BATCH')
    673                 "
    674             ).try_map(|r: PgRow| {
    675                 assert!(r.try_get(0)?);
    676                 Ok(())
    677             }).fetch_one(db).await.unwrap();
    678         }
    679 
    680         let now = Timestamp::now();
    681 
    682         // Submission retry status
    683         test(&db, async |batch_id| {
    684             batch_sub_failure(&db, batch_id, &now, "First failure")
    685                 .await
    686                 .unwrap();
    687             check_batch(batch_id, transient_failure, "First failure").await;
    688             batch_sub_failure(&db, batch_id, &now, "Second failure")
    689                 .await
    690                 .unwrap();
    691             check_batch(batch_id, transient_failure, "Second failure").await;
    692             batch_sub_success(&db, batch_id, &now, "ORDER")
    693                 .await
    694                 .unwrap();
    695             check_order("ORDER", pending, "").await;
    696             batch_sub_success(&db, batch_id, &now, "ORDER")
    697                 .await
    698                 .unwrap();
    699             check_order("ORDER", pending, "").await;
    700             order_step(&db, "ORDER", "step msg").await.unwrap();
    701             check_order("ORDER", pending, "step msg").await;
    702             order_step(&db, "ORDER", "success msg").await.unwrap();
    703             check_order("ORDER", pending, "success msg").await;
    704             order_success(&db, "ORDER").await.unwrap();
    705             check_order_tx("ORDER", success, "success msg", pending).await;
    706             order_step(&db, "ORDER", "late msg").await.unwrap();
    707             check_order_tx("ORDER", success, "success msg", pending).await;
    708         })
    709         .await;
    710 
    711         // Order step message on failure
    712         test(&db, async |batch_id| {
    713             batch_sub_success(&db, batch_id, &now, "ORDER")
    714                 .await
    715                 .unwrap();
    716             check_order("ORDER", pending, "").await;
    717             order_step(&db, "ORDER", "step msg").await.unwrap();
    718             check_order("ORDER", pending, "step msg").await;
    719             order_step(&db, "ORDER", "failure msg").await.unwrap();
    720             check_order("ORDER", pending, "failure msg").await;
    721             assert_eq!(
    722                 Some("failure msg"),
    723                 order_failure(&db, "ORDER")
    724                     .await
    725                     .unwrap()
    726                     .unwrap()
    727                     .1
    728                     .as_deref()
    729             );
    730             check_order("ORDER", permanent_failure, "failure msg").await;
    731             order_step(&db, "ORDER", "late msg").await.unwrap();
    732             check_order("ORDER", permanent_failure, "failure msg").await;
    733         })
    734         .await;
    735 
    736         // Payment & batch status
    737         test(&db, async |batch_id| {
    738             check_batch(batch_id, unsubmitted, "").await;
    739             batch_status_update(&db, "BATCH", pending, "progress")
    740                 .await
    741                 .unwrap();
    742             check_batch(batch_id, pending, "progress").await;
    743             tx_status_update(&db, "TX_SETTLED", "", success, "success")
    744                 .await
    745                 .unwrap();
    746             check_parts(
    747                 batch_id, pending, "progress", pending, "progress", success, "success",
    748             )
    749             .await;
    750             batch_status_update(&db, "BATCH", transient_failure, "waiting")
    751                 .await
    752                 .unwrap();
    753             check_parts(
    754                 batch_id,
    755                 transient_failure,
    756                 "waiting",
    757                 transient_failure,
    758                 "waiting",
    759                 success,
    760                 "success",
    761             )
    762             .await;
    763             tx_status_update(&db, "TX", "BATCH", permanent_failure, "failure")
    764                 .await
    765                 .unwrap();
    766             check_parts(
    767                 batch_id,
    768                 success,
    769                 "",
    770                 permanent_failure,
    771                 "failure",
    772                 success,
    773                 "success",
    774             )
    775             .await;
    776             tx_status_update(&db, "TX_SETTLED", "BATCH", permanent_failure, "late")
    777                 .await
    778                 .unwrap();
    779             check_parts(
    780                 batch_id,
    781                 success,
    782                 "",
    783                 permanent_failure,
    784                 "failure",
    785                 late_failure,
    786                 "late",
    787             )
    788             .await;
    789         })
    790         .await;
    791 
    792         // Registration
    793         test(&db, async |batch_id| {
    794             check_batch(batch_id, unsubmitted, "").await;
    795             register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX_SETTLED"))
    796                 .await
    797                 .unwrap();
    798             check_parts(batch_id, unsubmitted, "", unsubmitted, "", success, "").await;
    799             register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX").with_msg_id("BATCH"))
    800                 .await
    801                 .unwrap();
    802             check_parts(batch_id, success, "", success, "", success, "").await;
    803         })
    804         .await;
    805 
    806         // Transaction failure take over batch failures
    807         test(&db, async |batch_id| {
    808             check_batch(batch_id, unsubmitted, "").await;
    809             batch_status_update(&db, "BATCH", permanent_failure, "batch")
    810                 .await
    811                 .unwrap();
    812             check_parts(
    813                 batch_id,
    814                 permanent_failure,
    815                 "batch",
    816                 permanent_failure,
    817                 "batch",
    818                 permanent_failure,
    819                 "batch",
    820             )
    821             .await;
    822             tx_status_update(&db, "TX", "BATCH", permanent_failure, "tx")
    823                 .await
    824                 .unwrap();
    825             batch_status_update(&db, "BATCH", permanent_failure, "batch2")
    826                 .await
    827                 .unwrap();
    828             check_parts(
    829                 batch_id,
    830                 permanent_failure,
    831                 "batch",
    832                 permanent_failure,
    833                 "tx",
    834                 permanent_failure,
    835                 "batch",
    836             )
    837             .await;
    838         })
    839         .await;
    840 
    841         // Unknown order and batch
    842         batch_sub_success(&db, 42, &now, "ORDER_X").await.unwrap();
    843         batch_sub_failure(&db, 42, &now, "").await.unwrap();
    844         order_step(&db, "ORDER_X", "msg").await.unwrap();
    845         batch_status_update(&db, "BATCH_X", success, "")
    846             .await
    847             .unwrap();
    848         tx_status_update(&db, "TX_X", "BATCH_X", success, "msg")
    849             .await
    850             .unwrap();
    851         assert!(order_success(&db, "ORDER_X").await.unwrap().is_none());
    852         assert!(order_failure(&db, "ORDER_X").await.unwrap().is_none());
    853     }
    854 
    855     #[db_test]
    856     pub async fn initiated_submittables(db: PgPool) {
    857         let now = Timestamp::now();
    858         for i in 0..6 {
    859             assert!(matches!(
    860                 gen_initiate(&db, format!("PAY{i}"), "").await,
    861                 PaymentInitiationResult::Success(_)
    862             ));
    863             batch_initiated(&db, &now, &rand_ebics_id(), false)
    864                 .await
    865                 .unwrap();
    866         }
    867 
    868         let check_ids = async |ids: &[&str]| {
    869             assert_eq!(
    870                 ids,
    871                 initiated_submittable(&db, &CURR)
    872                     .await
    873                     .unwrap()
    874                     .iter()
    875                     .flat_map(|it| it.payments.iter().map(|it| it.e2e_id.as_str()))
    876                     .collect::<Vec<_>>()
    877             );
    878         };
    879         check_ids(&["PAY0", "PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await;
    880 
    881         // Check submitted not submitable
    882         batch_sub_success(&db, 1, &now, "ORDER1").await.unwrap();
    883         check_ids(&["PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await;
    884 
    885         // Check transient failure submitable last
    886         batch_sub_failure(&db, 2, &now, "Failure").await.unwrap();
    887         check_ids(&["PAY2", "PAY3", "PAY4", "PAY5", "PAY1"]).await;
    888 
    889         // Check persistent failure not submitable
    890         batch_sub_success(&db, 4, &now, "ORDER3").await.unwrap();
    891         order_failure(&db, "ORDER3").await.unwrap();
    892         check_ids(&["PAY2", "PAY4", "PAY5", "PAY1"]).await;
    893         batch_sub_success(&db, 5, &now, "ORDER4").await.unwrap();
    894         order_failure(&db, "ORDER4").await.unwrap();
    895         check_ids(&["PAY2", "PAY5", "PAY1"]).await;
    896 
    897         // Check rotation
    898         batch_sub_failure(&db, 3, &Timestamp::now(), "FAILURE")
    899             .await
    900             .unwrap();
    901         check_ids(&["PAY5", "PAY1", "PAY2"]).await;
    902         batch_sub_failure(&db, 6, &Timestamp::now(), "FAILURE")
    903             .await
    904             .unwrap();
    905         check_ids(&["PAY1", "PAY2", "PAY5"]).await;
    906         batch_sub_failure(&db, 2, &Timestamp::now(), "FAILURE")
    907             .await
    908             .unwrap();
    909         check_ids(&["PAY2", "PAY5", "PAY1"]).await;
    910     }
    911 }