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 }