db.rs (49903B)
1 /* 2 This file is part of TALER 3 Copyright (C) 2025, 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::fmt::Display; 18 19 use compact_str::CompactString; 20 use jiff::{Timestamp, civil::Date, tz::TimeZone}; 21 use serde::{Serialize, de::DeserializeOwned}; 22 use sqlx::{PgConnection, PgPool, QueryBuilder, Row, postgres::PgRow}; 23 use taler_api::{ 24 db::{BindHelper, TypeHelper, history, page}, 25 serialized, 26 subject::{IncomingKey, OutgoingSubject, fmt_out_subject}, 27 }; 28 use taler_common::{ 29 api::{ 30 HashCode, ShortHashCode, 31 params::{History, Page}, 32 prepared::{RegistrationRequest, Unregistration}, 33 revenue::RevenueIncomingBankTransaction, 34 wire::{ 35 IncomingBankTransaction, OutgoingBankTransaction, TransferListStatus, TransferState, 36 TransferStatus, 37 }, 38 }, 39 config::Config, 40 db::IncomingType, 41 types::{ 42 amount::{Amount, Decimal}, 43 payto::PaytoImpl as _, 44 }, 45 }; 46 use tokio::sync::watch::{Receiver, Sender}; 47 use url::Url; 48 49 use crate::{FullHuPayto, config::parse_db_cfg, constants::CURR, magnet_api::types::TxStatus}; 50 51 const SCHEMA: &str = "magnet_bank"; 52 53 pub async fn pool(cfg: &Config) -> anyhow::Result<PgPool> { 54 let db = parse_db_cfg(cfg)?; 55 let pool = taler_common::db::pool(db.cfg, SCHEMA).await?; 56 Ok(pool) 57 } 58 59 pub async fn dbinit(cfg: &Config, reset: bool) -> anyhow::Result<PgPool> { 60 let db_cfg = parse_db_cfg(cfg)?; 61 let pool = taler_common::db::pool(db_cfg.cfg, SCHEMA).await?; 62 let mut db = pool.acquire().await?; 63 taler_common::db::dbinit(&mut db, db_cfg.sql_dir.as_ref(), "magnet-bank", reset).await?; 64 Ok(pool) 65 } 66 67 pub async fn notification_listener( 68 pool: PgPool, 69 in_channel: Sender<i64>, 70 taler_in_channel: Sender<i64>, 71 out_channel: Sender<i64>, 72 taler_out_channel: Sender<i64>, 73 ) { 74 taler_api::notification::notification_listener!(&pool, 75 "tx_in" => (row_id: i64) { 76 in_channel.send_replace(row_id); 77 }, 78 "taler_in" => (row_id: i64) { 79 taler_in_channel.send_replace(row_id); 80 }, 81 "tx_out" => (row_id: i64) { 82 out_channel.send_replace(row_id); 83 }, 84 "taler_out" => (row_id: i64) { 85 taler_out_channel.send_replace(row_id); 86 } 87 ) 88 } 89 90 #[derive(Debug, Clone)] 91 pub struct TxIn { 92 pub code: u64, 93 pub amount: Amount, 94 pub subject: Box<str>, 95 pub debtor: FullHuPayto, 96 pub value_date: Date, 97 pub status: TxStatus, 98 } 99 100 impl Display for TxIn { 101 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { 102 let Self { 103 code, 104 amount, 105 subject, 106 debtor, 107 value_date, 108 status, 109 } = self; 110 write!( 111 f, 112 "{value_date} {code} {amount} ({} {}) {status:?} '{subject}'", 113 debtor.bban(), 114 debtor.name 115 ) 116 } 117 } 118 119 #[derive(Debug, Clone)] 120 pub struct TxOut { 121 pub code: u64, 122 pub amount: Amount, 123 pub subject: Box<str>, 124 pub creditor: FullHuPayto, 125 pub value_date: Date, 126 pub status: TxStatus, 127 } 128 129 impl Display for TxOut { 130 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { 131 let Self { 132 code, 133 amount, 134 subject, 135 creditor, 136 value_date, 137 status, 138 } = self; 139 write!( 140 f, 141 "{value_date} {code} {amount} ({} {}) {status:?} '{subject}'", 142 creditor.bban(), 143 creditor.name 144 ) 145 } 146 } 147 148 #[derive(Debug, PartialEq, Eq)] 149 pub struct Initiated { 150 pub id: u64, 151 pub amount: Amount, 152 pub subject: Box<str>, 153 pub creditor: FullHuPayto, 154 } 155 156 impl Display for Initiated { 157 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { 158 let Self { 159 id, 160 amount, 161 subject, 162 creditor, 163 } = self; 164 write!( 165 f, 166 "{id} {amount} ({} {}) '{subject}'", 167 creditor.bban(), 168 creditor.name 169 ) 170 } 171 } 172 173 #[derive(Debug, Clone)] 174 pub struct TxInAdmin { 175 pub amount: Amount, 176 pub subject: String, 177 pub debtor: FullHuPayto, 178 pub metadata: IncomingKey, 179 } 180 181 /// Lock the database for worker execution 182 pub async fn worker_lock(e: &mut PgConnection) -> sqlx::Result<bool> { 183 sqlx::query("SELECT pg_try_advisory_lock(42)") 184 .try_map(|r: PgRow| r.try_get(0)) 185 .fetch_one(e) 186 .await 187 } 188 189 #[derive(Debug, PartialEq, Eq)] 190 pub enum AddIncomingResult { 191 Success { 192 new: bool, 193 pending: bool, 194 row_id: u64, 195 valued_at: Date, 196 }, 197 ReservePubReuse, 198 UnknownMapping, 199 MappingReuse, 200 } 201 202 pub async fn register_tx_in_admin( 203 db: &PgPool, 204 tx: &TxInAdmin, 205 now: &Timestamp, 206 ) -> sqlx::Result<AddIncomingResult> { 207 serialized!( 208 sqlx::query( 209 " 210 SELECT out_reserve_pub_reuse, out_mapping_reuse, out_unknown_mapping, out_tx_row_id, out_valued_at, out_new, out_pending 211 FROM register_tx_in(NULL, $1, $2, $3, $4, $5, $6, $7, $5) 212 ", 213 ) 214 .bind(tx.amount) 215 .bind(&tx.subject) 216 .bind(tx.debtor.iban()) 217 .bind(&tx.debtor.name) 218 .bind_date(&now.to_zoned(TimeZone::UTC).date()) 219 .bind(tx.metadata.ty) 220 .bind(tx.metadata.key) 221 .try_map(|r: PgRow| { 222 Ok(if r.try_get_flag(0)? { 223 AddIncomingResult::ReservePubReuse 224 } else if r.try_get_flag(1)? { 225 AddIncomingResult::MappingReuse 226 } else if r.try_get_flag(2)? { 227 AddIncomingResult::UnknownMapping 228 } else { 229 AddIncomingResult::Success { 230 row_id: r.try_get_u64(3)?, 231 valued_at: r.try_get_date(4)?, 232 new: r.try_get(5)?, 233 pending: r.try_get(6)? 234 } 235 }) 236 }) 237 .fetch_one(db) 238 ) 239 } 240 241 pub async fn register_tx_in( 242 db: &mut PgConnection, 243 tx: &TxIn, 244 subject: &Option<IncomingKey>, 245 now: &Timestamp, 246 ) -> sqlx::Result<AddIncomingResult> { 247 serialized!( 248 sqlx::query( 249 " 250 SELECT out_reserve_pub_reuse, out_mapping_reuse, out_unknown_mapping, out_tx_row_id, out_valued_at, out_new, out_pending 251 FROM register_tx_in($1, $2, $3, $4, $5, $6, $7, $8, $9) 252 ", 253 ) 254 .bind(tx.code as i64) 255 .bind(tx.amount) 256 .bind(&tx.subject) 257 .bind(tx.debtor.iban()) 258 .bind(&tx.debtor.name) 259 .bind_date(&tx.value_date) 260 .bind(subject.as_ref().map(|it| it.ty)) 261 .bind(subject.as_ref().map(|it| it.key)) 262 .bind_timestamp(now) 263 .try_map(|r: PgRow| { 264 Ok(if r.try_get_flag(0)? { 265 AddIncomingResult::ReservePubReuse 266 } else if r.try_get_flag(1)? { 267 AddIncomingResult::MappingReuse 268 } else if r.try_get_flag(2)? { 269 AddIncomingResult::UnknownMapping 270 } else { 271 AddIncomingResult::Success { 272 row_id: r.try_get_u64(3)?, 273 valued_at: r.try_get_date(4)?, 274 new: r.try_get(5)?, 275 pending: r.try_get(6)? 276 } 277 }) 278 }) 279 .fetch_one(&mut *db) 280 ) 281 } 282 283 #[derive(Debug)] 284 pub enum TxOutKind { 285 Simple, 286 Bounce(u32), 287 Talerable(OutgoingSubject), 288 } 289 290 #[derive(Debug, Clone, Copy, PartialEq, Eq, sqlx::Type)] 291 #[allow(non_camel_case_types)] 292 #[sqlx(type_name = "register_result")] 293 pub enum RegisterResult { 294 /// Already registered 295 idempotent, 296 /// Initiated transaction 297 known, 298 /// Recovered unknown outgoing transaction 299 recovered, 300 } 301 302 #[derive(Debug, PartialEq, Eq)] 303 pub struct AddOutgoingResult { 304 pub result: RegisterResult, 305 pub row_id: u64, 306 } 307 308 pub async fn register_tx_out( 309 db: &mut PgConnection, 310 tx: &TxOut, 311 kind: &TxOutKind, 312 now: &Timestamp, 313 ) -> sqlx::Result<AddOutgoingResult> { 314 serialized!({ 315 let query = sqlx::query( 316 " 317 SELECT out_result, out_tx_row_id 318 FROM register_tx_out($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) 319 ", 320 ) 321 .bind(tx.code as i64) 322 .bind(tx.amount) 323 .bind(&tx.subject) 324 .bind(tx.creditor.iban()) 325 .bind(&tx.creditor.name) 326 .bind_date(&tx.value_date); 327 match kind { 328 TxOutKind::Simple => query 329 .bind(None::<&[u8]>) 330 .bind(None::<&str>) 331 .bind(None::<&str>) 332 .bind(None::<i64>), 333 TxOutKind::Bounce(bounced) => query 334 .bind(None::<&[u8]>) 335 .bind(None::<&str>) 336 .bind(None::<&str>) 337 .bind(*bounced as i64), 338 TxOutKind::Talerable(subject) => query 339 .bind(subject.wtid) 340 .bind(subject.exchange_base_url.as_str()) 341 .bind(&subject.metadata) 342 .bind(None::<i64>), 343 } 344 .bind_timestamp(now) 345 .try_map(|r: PgRow| { 346 Ok(AddOutgoingResult { 347 result: r.try_get(0)?, 348 row_id: r.try_get_u64(1)?, 349 }) 350 }) 351 .fetch_one(&mut *db) 352 }) 353 } 354 355 #[derive(Debug, PartialEq, Eq)] 356 pub struct OutFailureResult { 357 pub initiated_id: Option<u64>, 358 pub new: bool, 359 } 360 361 pub async fn register_tx_out_failure( 362 db: &mut PgConnection, 363 code: u64, 364 bounced: Option<u32>, 365 now: &Timestamp, 366 ) -> sqlx::Result<OutFailureResult> { 367 serialized!( 368 sqlx::query( 369 " 370 SELECT out_new, out_initiated_id 371 FROM register_tx_out_failure($1, $2, $3) 372 ", 373 ) 374 .bind(code as i64) 375 .bind(bounced.map(|i| i as i32)) 376 .bind_timestamp(now) 377 .try_map(|r: PgRow| { 378 Ok(OutFailureResult { 379 new: r.try_get(0)?, 380 initiated_id: r.try_get::<Option<i64>, _>(1)?.map(|i| i as u64), 381 }) 382 }) 383 .fetch_one(&mut *db) 384 ) 385 } 386 387 #[derive(Debug, PartialEq, Eq)] 388 pub enum TransferResult { 389 Success { id: u64, initiated_at: Timestamp }, 390 RequestUidReuse, 391 WtidReuse, 392 } 393 394 #[derive(Debug, Clone)] 395 pub struct Transfer { 396 pub request_uid: HashCode, 397 pub amount: Decimal, 398 pub exchange_base_url: Url, 399 pub metadata: Option<CompactString>, 400 pub wtid: ShortHashCode, 401 pub creditor: FullHuPayto, 402 } 403 404 pub async fn make_transfer( 405 db: &PgPool, 406 tx: &Transfer, 407 now: &Timestamp, 408 ) -> sqlx::Result<TransferResult> { 409 let subject = fmt_out_subject(&tx.wtid, &tx.exchange_base_url, tx.metadata.as_deref()); 410 serialized!( 411 sqlx::query( 412 " 413 SELECT out_request_uid_reuse, out_wtid_reuse, out_initiated_row_id, out_initiated_at 414 FROM taler_transfer($1, $2, $3, $4, $5, $6, $7, $8, $9) 415 ", 416 ) 417 .bind(tx.request_uid) 418 .bind(tx.wtid) 419 .bind(&subject) 420 .bind(tx.amount) 421 .bind(tx.exchange_base_url.as_str()) 422 .bind(&tx.metadata) 423 .bind(tx.creditor.iban()) 424 .bind(&tx.creditor.name) 425 .bind_timestamp(now) 426 .try_map(|r: PgRow| { 427 Ok(if r.try_get_flag(0)? { 428 TransferResult::RequestUidReuse 429 } else if r.try_get_flag(1)? { 430 TransferResult::WtidReuse 431 } else { 432 TransferResult::Success { 433 id: r.try_get_u64(2)?, 434 initiated_at: r.try_get_timestamp(3)?, 435 } 436 }) 437 }) 438 .fetch_one(db) 439 ) 440 } 441 442 #[derive(Debug, PartialEq, Eq)] 443 pub struct BounceResult { 444 pub tx_id: u64, 445 pub tx_new: bool, 446 pub bounce_id: u64, 447 pub bounce_new: bool, 448 } 449 450 pub async fn register_bounce_tx_in( 451 db: &mut PgConnection, 452 tx: &TxIn, 453 reason: &str, 454 now: &Timestamp, 455 ) -> sqlx::Result<BounceResult> { 456 serialized!( 457 sqlx::query( 458 " 459 SELECT out_tx_row_id, out_tx_new, out_bounce_row_id, out_bounce_new 460 FROM register_bounce_tx_in($1, $2, $3, $4, $5, $6, $7, $8) 461 ", 462 ) 463 .bind(tx.code as i64) 464 .bind(tx.amount) 465 .bind(&tx.subject) 466 .bind(tx.debtor.iban()) 467 .bind(&tx.debtor.name) 468 .bind_date(&tx.value_date) 469 .bind(reason) 470 .bind_timestamp(now) 471 .try_map(|r: PgRow| { 472 Ok(BounceResult { 473 tx_id: r.try_get_u64(0)?, 474 tx_new: r.try_get(1)?, 475 bounce_id: r.try_get_u64(2)?, 476 bounce_new: r.try_get(3)?, 477 }) 478 }) 479 .fetch_one(&mut *db) 480 ) 481 } 482 483 pub async fn transfer_page( 484 db: &PgPool, 485 status: &Option<TransferState>, 486 params: &Page, 487 ) -> sqlx::Result<Vec<TransferListStatus>> { 488 page( 489 db, 490 params, 491 "initiated_id", 492 || { 493 let mut builder = QueryBuilder::new( 494 " 495 SELECT 496 initiated_id, 497 status, 498 amount, 499 credit_account, 500 credit_name, 501 initiated_at 502 FROM transfer 503 JOIN initiated USING (initiated_id) 504 WHERE 505 ", 506 ); 507 if let Some(status) = status { 508 builder.push(" status = ").push_bind(status).push(" AND "); 509 } 510 builder 511 }, 512 |r: PgRow| { 513 Ok(TransferListStatus { 514 row_id: r.try_get_u64(0)?, 515 status: r.try_get(1)?, 516 amount: r.try_get_amount(2, &CURR)?, 517 credit_account: r.try_get_iban(3)?.as_full_uri(r.try_get(4)?), 518 timestamp: r.try_get_timestamp(5)?.into(), 519 }) 520 }, 521 ) 522 .await 523 } 524 525 pub async fn outgoing_history( 526 db: &PgPool, 527 params: &History, 528 listen: impl FnOnce() -> Receiver<i64>, 529 ) -> sqlx::Result<Vec<OutgoingBankTransaction>> { 530 history( 531 db, 532 "tx_out_id", 533 params, 534 listen, 535 || { 536 QueryBuilder::new( 537 " 538 SELECT 539 tx_out_id, 540 amount, 541 credit_account, 542 credit_name, 543 valued_at, 544 exchange_base_url, 545 metadata, 546 wtid 547 FROM taler_out 548 JOIN tx_out USING (tx_out_id) 549 WHERE 550 ", 551 ) 552 }, 553 |r: PgRow| { 554 Ok(OutgoingBankTransaction { 555 row_id: r.try_get_u64(0)?, 556 amount: r.try_get_amount(1, &CURR)?, 557 debit_fee: None, 558 credit_account: r.try_get_iban(2)?.as_full_uri(r.try_get(3)?), 559 date: r.try_get_timestamp(4)?.into(), 560 exchange_base_url: r.try_get_url(5)?, 561 metadata: r.try_get(6)?, 562 wtid: r.try_get(7)?, 563 }) 564 }, 565 ) 566 .await 567 } 568 569 pub async fn incoming_history( 570 db: &PgPool, 571 params: &History, 572 listen: impl FnOnce() -> Receiver<i64>, 573 ) -> sqlx::Result<Vec<IncomingBankTransaction>> { 574 history( 575 db, 576 "tx_in_id", 577 params, 578 listen, 579 || { 580 QueryBuilder::new( 581 " 582 SELECT 583 type, 584 tx_in_id, 585 amount, 586 debit_account, 587 debit_name, 588 valued_at, 589 metadata, 590 authorization_pub, 591 authorization_sig 592 FROM taler_in 593 JOIN tx_in USING (tx_in_id) 594 WHERE 595 ", 596 ) 597 }, 598 |r: PgRow| { 599 Ok(match r.try_get(0)? { 600 IncomingType::reserve => IncomingBankTransaction::Reserve { 601 row_id: r.try_get_u64(1)?, 602 amount: r.try_get_amount(2, &CURR)?, 603 credit_fee: None, 604 debit_account: r.try_get_iban(3)?.as_full_uri(r.try_get(4)?), 605 date: r.try_get_timestamp(5)?.into(), 606 reserve_pub: r.try_get(6)?, 607 authorization_pub: r.try_get(7)?, 608 authorization_sig: r.try_get(8)?, 609 }, 610 IncomingType::kyc => IncomingBankTransaction::Kyc { 611 row_id: r.try_get_u64(1)?, 612 amount: r.try_get_amount(2, &CURR)?, 613 credit_fee: None, 614 debit_account: r.try_get_iban(3)?.as_full_uri(r.try_get(4)?), 615 date: r.try_get_timestamp(5)?.into(), 616 account_pub: r.try_get(6)?, 617 authorization_pub: r.try_get(7)?, 618 authorization_sig: r.try_get(8)?, 619 }, 620 IncomingType::map => unimplemented!("MAP are never listed in the history"), 621 }) 622 }, 623 ) 624 .await 625 } 626 627 pub async fn revenue_history( 628 db: &PgPool, 629 params: &History, 630 listen: impl FnOnce() -> Receiver<i64>, 631 ) -> sqlx::Result<Vec<RevenueIncomingBankTransaction>> { 632 history( 633 db, 634 "tx_in_id", 635 params, 636 listen, 637 || { 638 QueryBuilder::new( 639 " 640 SELECT 641 tx_in_id, 642 valued_at, 643 amount, 644 debit_account, 645 debit_name, 646 subject 647 FROM tx_in 648 WHERE 649 ", 650 ) 651 }, 652 |r: PgRow| { 653 Ok(RevenueIncomingBankTransaction { 654 row_id: r.try_get_u64(0)?, 655 date: r.try_get_timestamp(1)?.into(), 656 amount: r.try_get_amount(2, &CURR)?, 657 credit_fee: None, 658 debit_account: r.try_get_iban(3)?.as_full_uri(r.try_get(4)?), 659 subject: r.try_get(5)?, 660 }) 661 }, 662 ) 663 .await 664 } 665 666 pub async fn transfer_by_id(db: &PgPool, id: u64) -> sqlx::Result<Option<TransferStatus>> { 667 serialized!( 668 sqlx::query( 669 " 670 SELECT 671 status, 672 status_msg, 673 amount, 674 exchange_base_url, 675 metadata, 676 wtid, 677 credit_account, 678 credit_name, 679 initiated_at 680 FROM transfer 681 JOIN initiated USING (initiated_id) 682 WHERE initiated_id = $1 683 ", 684 ) 685 .bind(id as i64) 686 .try_map(|r: PgRow| { 687 Ok(TransferStatus { 688 status: r.try_get(0)?, 689 status_msg: r.try_get(1)?, 690 amount: r.try_get_amount(2, &CURR)?, 691 exchange_base_url: r.try_get(3)?, 692 metadata: r.try_get(4)?, 693 wtid: r.try_get(5)?, 694 credit_account: r.try_get_iban(6)?.as_full_uri(r.try_get(7)?), 695 timestamp: r.try_get_timestamp(8)?.into(), 696 }) 697 }) 698 .fetch_optional(db) 699 ) 700 } 701 702 /** Get a batch of pending initiated transactions not attempted since [start] */ 703 pub async fn pending_batch( 704 db: &mut PgConnection, 705 start: &Timestamp, 706 ) -> sqlx::Result<Vec<Initiated>> { 707 serialized!( 708 sqlx::query( 709 " 710 SELECT initiated_id, amount, subject, credit_account, credit_name 711 FROM initiated 712 WHERE magnet_code IS NULL 713 AND status='pending' 714 AND (last_submitted IS NULL OR last_submitted < $1) 715 LIMIT 100 716 ", 717 ) 718 .bind_timestamp(start) 719 .try_map(|r: PgRow| { 720 Ok(Initiated { 721 id: r.try_get_u64(0)?, 722 amount: r.try_get_amount(1, &CURR)?, 723 subject: r.try_get(2)?, 724 creditor: FullHuPayto::new(r.try_get_parse(3)?, r.try_get(4)?), 725 }) 726 }) 727 .fetch_all(&mut *db) 728 ) 729 } 730 731 /** Get an initiated transaction matching the given magnet [code] */ 732 pub async fn initiated_by_code( 733 db: &mut PgConnection, 734 code: u64, 735 ) -> sqlx::Result<Option<Initiated>> { 736 serialized!( 737 sqlx::query( 738 " 739 SELECT initiated_id, amount, subject, credit_account, credit_name 740 FROM initiated 741 WHERE magnet_code IS $1 742 ", 743 ) 744 .bind(code as i64) 745 .try_map(|r: PgRow| { 746 Ok(Initiated { 747 id: r.try_get_u64(0)?, 748 amount: r.try_get_amount(1, &CURR)?, 749 subject: r.try_get(2)?, 750 creditor: FullHuPayto::new(r.try_get_parse(3)?, r.try_get(4)?), 751 }) 752 }) 753 .fetch_optional(&mut *db) 754 ) 755 } 756 757 /** Update status of a successful submitted initiated transaction */ 758 pub async fn initiated_submit_success( 759 db: &mut PgConnection, 760 id: u64, 761 timestamp: &Timestamp, 762 magnet_code: u64, 763 ) -> sqlx::Result<()> { 764 serialized!( 765 sqlx::query( 766 " 767 UPDATE initiated 768 SET status='pending', submission_counter=submission_counter+1, last_submitted=$1, magnet_code=$2 769 WHERE initiated_id=$3 770 " 771 ).bind_timestamp(timestamp) 772 .bind(magnet_code as i64) 773 .bind(id as i64) 774 .execute(&mut *db) 775 )?; 776 Ok(()) 777 } 778 779 /** Update status of a permanently failed initiated transaction */ 780 pub async fn initiated_submit_permanent_failure( 781 db: &mut PgConnection, 782 id: u64, 783 timestamp: &Timestamp, 784 msg: &str, 785 ) -> sqlx::Result<()> { 786 serialized!( 787 sqlx::query( 788 " 789 UPDATE initiated 790 SET status='permanent_failure', status_msg=$2 791 WHERE initiated_id=$3 792 ", 793 ) 794 .bind_timestamp(timestamp) 795 .bind(msg) 796 .bind(id as i64) 797 .execute(&mut *db) 798 )?; 799 Ok(()) 800 } 801 802 /** Check if an initiated transaction exist for a magnet code */ 803 pub async fn initiated_exists_for_code( 804 db: &mut PgConnection, 805 code: u64, 806 ) -> sqlx::Result<Option<u64>> { 807 serialized!( 808 sqlx::query("SELECT initiated_id FROM initiated WHERE magnet_code=$1") 809 .bind(code as i64) 810 .try_map(|r| Ok(r.try_get::<i64, _>(0)? as u64)) 811 .fetch_optional(&mut *db) 812 ) 813 } 814 815 /** Get JSON value from KV table */ 816 pub async fn kv_get<T: DeserializeOwned + Unpin + Send>( 817 db: &mut PgConnection, 818 key: &str, 819 ) -> sqlx::Result<Option<T>> { 820 serialized!( 821 sqlx::query("SELECT value FROM kv WHERE key=$1") 822 .bind(key) 823 .try_map(|r| Ok(r.try_get::<sqlx::types::Json<T>, _>(0)?.0)) 824 .fetch_optional(&mut *db) 825 ) 826 } 827 828 /** Set JSON value in KV table */ 829 pub async fn kv_set<T: Serialize>(db: &mut PgConnection, key: &str, value: &T) -> sqlx::Result<()> { 830 serialized!( 831 sqlx::query("INSERT INTO kv (key, value) VALUES ($1, $2) ON CONFLICT (key) DO UPDATE SET value=EXCLUDED.value") 832 .bind(key) 833 .bind(sqlx::types::Json(value)) 834 .execute(&mut *db) 835 )?; 836 Ok(()) 837 } 838 839 pub enum RegistrationResult { 840 Success, 841 ReservePubReuse, 842 } 843 844 pub async fn transfer_register( 845 db: &PgPool, 846 req: &RegistrationRequest, 847 ) -> sqlx::Result<RegistrationResult> { 848 let ty: IncomingType = req.r#type.into(); 849 serialized!( 850 sqlx::query( 851 "SELECT out_reserve_pub_reuse FROM register_prepared_transfers($1,$2,$3,$4,$5,$6)" 852 ) 853 .bind(ty) 854 .bind(req.account_pub) 855 .bind(req.authorization_pub) 856 .bind(req.authorization_sig) 857 .bind(req.recurrent) 858 .bind_timestamp(&Timestamp::now()) 859 .try_map(|r: PgRow| { 860 Ok(if r.try_get_flag("out_reserve_pub_reuse")? { 861 RegistrationResult::ReservePubReuse 862 } else { 863 RegistrationResult::Success 864 }) 865 }) 866 .fetch_one(db) 867 ) 868 } 869 870 pub async fn transfer_unregister(db: &PgPool, req: &Unregistration) -> sqlx::Result<bool> { 871 serialized!( 872 sqlx::query("SELECT out_found FROM delete_prepared_transfers($1,$2)") 873 .bind(req.authorization_pub) 874 .bind_timestamp(&Timestamp::now()) 875 .try_map(|r: PgRow| r.try_get_flag("out_found")) 876 .fetch_one(db) 877 ) 878 } 879 880 #[cfg(test)] 881 mod test { 882 use jiff::{Span, Timestamp, Zoned, tz::TimeZone}; 883 use serde_json::json; 884 use sqlx::{PgConnection, PgPool, postgres::PgRow}; 885 use taler_api::{ 886 db::TypeHelper, 887 notification::dummy_listen, 888 subject::{IncomingKey, OutgoingSubject}, 889 }; 890 use taler_common::{ 891 api::{ 892 EddsaPublicKey, HashCode, ShortHashCode, 893 params::{History, Page}, 894 }, 895 types::{ 896 amount::{amount, decimal}, 897 url, 898 utils::now_sql_stable_ts, 899 }, 900 }; 901 use taler_macros::db_test; 902 903 use super::TxInAdmin; 904 use crate::{ 905 db::{ 906 self, AddIncomingResult, AddOutgoingResult, BounceResult, Initiated, OutFailureResult, 907 TransferResult, TxIn, TxOut, TxOutKind, kv_get, kv_set, make_transfer, 908 register_bounce_tx_in, register_tx_in, register_tx_in_admin, register_tx_out, 909 }, 910 magnet_api::types::TxStatus, 911 magnet_payto, 912 }; 913 914 #[db_test] 915 async fn kv(mut db: PgConnection) { 916 let value = json!({ 917 "name": "Mr Smith", 918 "no way": 32 919 }); 920 921 assert_eq!( 922 kv_get::<serde_json::Value>(&mut db, "value").await.unwrap(), 923 None 924 ); 925 kv_set(&mut db, "value", &value).await.unwrap(); 926 kv_set(&mut db, "value", &value).await.unwrap(); 927 assert_eq!( 928 kv_get::<serde_json::Value>(&mut db, "value").await.unwrap(), 929 Some(value) 930 ); 931 } 932 933 #[db_test] 934 async fn tx_in(mut db: PgConnection, pool: PgPool) { 935 let mut routine = async |first: &Option<IncomingKey>, second: &Option<IncomingKey>| { 936 let (id, code) = 937 sqlx::query("SELECT count(*) + 1, COALESCE(max(magnet_code), 0) + 20 FROM tx_in") 938 .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get_u64(1)?))) 939 .fetch_one(&mut db) 940 .await 941 .unwrap(); 942 let now = now_sql_stable_ts(); 943 let date = Zoned::now().date(); 944 let later = date.tomorrow().unwrap(); 945 let tx = TxIn { 946 code, 947 amount: amount("EUR:10"), 948 subject: "subject".into(), 949 debtor: magnet_payto( 950 "payto://iban/HU30162000031000163100000000?receiver-name=name", 951 ), 952 value_date: date, 953 status: TxStatus::Completed, 954 }; 955 // Insert 956 assert_eq!( 957 register_tx_in(&mut db, &tx, first, &now) 958 .await 959 .expect("register tx in"), 960 AddIncomingResult::Success { 961 new: true, 962 pending: false, 963 row_id: id, 964 valued_at: date 965 } 966 ); 967 // Idempotent 968 assert_eq!( 969 register_tx_in( 970 &mut db, 971 &TxIn { 972 value_date: later, 973 ..tx.clone() 974 }, 975 first, 976 &now 977 ) 978 .await 979 .expect("register tx in"), 980 AddIncomingResult::Success { 981 new: false, 982 pending: false, 983 row_id: id, 984 valued_at: date 985 } 986 ); 987 // Many 988 assert_eq!( 989 register_tx_in( 990 &mut db, 991 &TxIn { 992 code: code + 1, 993 value_date: later, 994 ..tx 995 }, 996 second, 997 &now 998 ) 999 .await 1000 .expect("register tx in"), 1001 AddIncomingResult::Success { 1002 new: true, 1003 pending: false, 1004 row_id: id + 1, 1005 valued_at: later 1006 } 1007 ); 1008 }; 1009 1010 // Empty db 1011 assert_eq!( 1012 db::revenue_history(&pool, &History::default(), dummy_listen) 1013 .await 1014 .unwrap(), 1015 Vec::new() 1016 ); 1017 assert_eq!( 1018 db::incoming_history(&pool, &History::default(), dummy_listen) 1019 .await 1020 .unwrap(), 1021 Vec::new() 1022 ); 1023 1024 // Regular transaction 1025 routine(&None, &None).await; 1026 1027 // Reserve transaction 1028 routine( 1029 &Some(IncomingKey::reserve(EddsaPublicKey::rand())), 1030 &Some(IncomingKey::reserve(EddsaPublicKey::rand())), 1031 ) 1032 .await; 1033 1034 // Kyc transaction 1035 routine( 1036 &Some(IncomingKey::kyc(EddsaPublicKey::rand())), 1037 &Some(IncomingKey::kyc(EddsaPublicKey::rand())), 1038 ) 1039 .await; 1040 1041 // History 1042 assert_eq!( 1043 db::revenue_history(&pool, &History::default(), dummy_listen) 1044 .await 1045 .unwrap() 1046 .len(), 1047 6 1048 ); 1049 assert_eq!( 1050 db::incoming_history(&pool, &History::default(), dummy_listen) 1051 .await 1052 .unwrap() 1053 .len(), 1054 4 1055 ); 1056 } 1057 1058 #[db_test] 1059 async fn tx_in_admin(pool: PgPool) { 1060 // Empty db 1061 assert_eq!( 1062 db::incoming_history(&pool, &History::default(), dummy_listen) 1063 .await 1064 .unwrap(), 1065 Vec::new() 1066 ); 1067 1068 let now: Timestamp = "2026-09-04T23:00:00Z".parse().unwrap(); 1069 let later = now + Span::new().hours(2); 1070 let date = now.to_zoned(TimeZone::UTC).date(); 1071 let later_date = later.to_zoned(TimeZone::UTC).date(); 1072 let tx = TxInAdmin { 1073 amount: amount("EUR:10"), 1074 subject: "subject".to_owned(), 1075 debtor: magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"), 1076 metadata: IncomingKey::reserve(EddsaPublicKey::rand()), 1077 }; 1078 // Insert 1079 assert_eq!( 1080 register_tx_in_admin(&pool, &tx, &now) 1081 .await 1082 .expect("register tx in"), 1083 AddIncomingResult::Success { 1084 new: true, 1085 pending: false, 1086 row_id: 1, 1087 valued_at: date 1088 } 1089 ); 1090 // Many 1091 assert_eq!( 1092 register_tx_in_admin( 1093 &pool, 1094 &TxInAdmin { 1095 subject: "Other".to_owned(), 1096 metadata: IncomingKey::reserve(EddsaPublicKey::rand()), 1097 ..tx.clone() 1098 }, 1099 &later 1100 ) 1101 .await 1102 .expect("register tx in"), 1103 AddIncomingResult::Success { 1104 new: true, 1105 pending: false, 1106 row_id: 2, 1107 valued_at: later_date 1108 } 1109 ); 1110 1111 // History 1112 assert_eq!( 1113 db::incoming_history(&pool, &History::default(), dummy_listen) 1114 .await 1115 .unwrap() 1116 .len(), 1117 2 1118 ); 1119 } 1120 1121 #[db_test] 1122 async fn tx_out(mut db: PgConnection, pool: PgPool) { 1123 let mut routine = async |first: &TxOutKind, second: &TxOutKind| { 1124 let (id, code) = 1125 sqlx::query("SELECT count(*) + 1, COALESCE(max(magnet_code), 0) + 20 FROM tx_out") 1126 .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get_u64(1)?))) 1127 .fetch_one(&mut db) 1128 .await 1129 .unwrap(); 1130 let now = now_sql_stable_ts(); 1131 let date = Zoned::now().date(); 1132 let later = date.tomorrow().unwrap(); 1133 let tx = TxOut { 1134 code, 1135 amount: amount("HUF:10"), 1136 subject: "subject".into(), 1137 creditor: magnet_payto( 1138 "payto://iban/HU30162000031000163100000000?receiver-name=name", 1139 ), 1140 value_date: date, 1141 status: TxStatus::Completed, 1142 }; 1143 // TODO revert to assert_matches! once rust-lang#82775 is stable 1144 assert!(matches!( 1145 make_transfer( 1146 &pool, 1147 &db::Transfer { 1148 request_uid: HashCode::rand(), 1149 amount: decimal("10"), 1150 exchange_base_url: url("https://exchange.test.com/"), 1151 metadata: None, 1152 wtid: ShortHashCode::rand(), 1153 creditor: tx.creditor.clone() 1154 }, 1155 &now 1156 ) 1157 .await 1158 .unwrap(), 1159 TransferResult::Success { .. } 1160 )); 1161 db::initiated_submit_success(&mut db, 1, &Timestamp::now(), tx.code) 1162 .await 1163 .expect("status success"); 1164 1165 // Insert 1166 assert_eq!( 1167 register_tx_out(&mut db, &tx, first, &now) 1168 .await 1169 .expect("register tx out"), 1170 AddOutgoingResult { 1171 result: db::RegisterResult::known, 1172 row_id: id, 1173 } 1174 ); 1175 // Idempotent 1176 assert_eq!( 1177 register_tx_out( 1178 &mut db, 1179 &TxOut { 1180 value_date: later, 1181 ..tx.clone() 1182 }, 1183 first, 1184 &now 1185 ) 1186 .await 1187 .expect("register tx out"), 1188 AddOutgoingResult { 1189 result: db::RegisterResult::idempotent, 1190 row_id: id, 1191 } 1192 ); 1193 // Recovered 1194 assert_eq!( 1195 register_tx_out( 1196 &mut db, 1197 &TxOut { 1198 code: code + 1, 1199 value_date: later, 1200 ..tx.clone() 1201 }, 1202 second, 1203 &now 1204 ) 1205 .await 1206 .expect("register tx out"), 1207 AddOutgoingResult { 1208 result: db::RegisterResult::recovered, 1209 row_id: id + 1, 1210 } 1211 ); 1212 }; 1213 1214 // Empty db 1215 assert_eq!( 1216 db::outgoing_history(&pool, &History::default(), dummy_listen) 1217 .await 1218 .unwrap(), 1219 Vec::new() 1220 ); 1221 1222 // Regular transaction 1223 routine(&TxOutKind::Simple, &TxOutKind::Simple).await; 1224 1225 // Talerable transaction 1226 routine( 1227 &TxOutKind::Talerable(OutgoingSubject::rand()), 1228 &TxOutKind::Talerable(OutgoingSubject::rand()), 1229 ) 1230 .await; 1231 1232 // Bounced transaction 1233 routine(&TxOutKind::Bounce(21), &TxOutKind::Bounce(42)).await; 1234 1235 // History 1236 assert_eq!( 1237 db::outgoing_history(&pool, &History::default(), dummy_listen) 1238 .await 1239 .unwrap() 1240 .len(), 1241 2 1242 ); 1243 } 1244 1245 #[db_test] 1246 async fn tx_out_failure(mut db: PgConnection, pool: PgPool) { 1247 let now = now_sql_stable_ts(); 1248 1249 // Unknown 1250 assert_eq!( 1251 db::register_tx_out_failure(&mut db, 42, None, &now) 1252 .await 1253 .unwrap(), 1254 OutFailureResult { 1255 initiated_id: None, 1256 new: false 1257 } 1258 ); 1259 assert_eq!( 1260 db::register_tx_out_failure(&mut db, 42, Some(12), &now) 1261 .await 1262 .unwrap(), 1263 OutFailureResult { 1264 initiated_id: None, 1265 new: false 1266 } 1267 ); 1268 1269 // Initiated 1270 let req = db::Transfer { 1271 request_uid: HashCode::rand(), 1272 amount: decimal("10"), 1273 exchange_base_url: url("https://exchange.test.com/"), 1274 metadata: None, 1275 wtid: ShortHashCode::rand(), 1276 creditor: magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"), 1277 }; 1278 let payto = magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"); 1279 assert_eq!( 1280 make_transfer(&pool, &req, &now).await.unwrap(), 1281 TransferResult::Success { 1282 id: 1, 1283 initiated_at: now 1284 } 1285 ); 1286 db::initiated_submit_success(&mut db, 1, &Timestamp::now(), 34) 1287 .await 1288 .expect("status success"); 1289 assert_eq!( 1290 db::register_tx_out_failure(&mut db, 34, None, &now) 1291 .await 1292 .unwrap(), 1293 OutFailureResult { 1294 initiated_id: Some(1), 1295 new: true 1296 } 1297 ); 1298 assert_eq!( 1299 db::register_tx_out_failure(&mut db, 34, None, &now) 1300 .await 1301 .unwrap(), 1302 OutFailureResult { 1303 initiated_id: Some(1), 1304 new: false 1305 } 1306 ); 1307 1308 // Recovered bounce 1309 let tx = TxIn { 1310 code: 12, 1311 amount: amount("HUF:11"), 1312 subject: "malformed transaction".into(), 1313 debtor: payto, 1314 value_date: Zoned::now().date(), 1315 status: TxStatus::Completed, 1316 }; 1317 assert_eq!( 1318 db::register_bounce_tx_in(&mut db, &tx, "no reason", &now) 1319 .await 1320 .unwrap(), 1321 BounceResult { 1322 tx_id: 1, 1323 tx_new: true, 1324 bounce_id: 2, 1325 bounce_new: true 1326 } 1327 ); 1328 assert_eq!( 1329 db::register_tx_out_failure(&mut db, 10, Some(12), &now) 1330 .await 1331 .unwrap(), 1332 OutFailureResult { 1333 initiated_id: Some(2), 1334 new: true 1335 } 1336 ); 1337 assert_eq!( 1338 db::register_tx_out_failure(&mut db, 10, Some(12), &now) 1339 .await 1340 .unwrap(), 1341 OutFailureResult { 1342 initiated_id: Some(2), 1343 new: false 1344 } 1345 ); 1346 } 1347 1348 #[db_test] 1349 async fn transfer(pool: PgPool) { 1350 // Empty db 1351 assert_eq!(db::transfer_by_id(&pool, 0).await.unwrap(), None); 1352 assert_eq!( 1353 db::transfer_page(&pool, &None, &Page::default()) 1354 .await 1355 .unwrap(), 1356 Vec::new() 1357 ); 1358 1359 let req = db::Transfer { 1360 request_uid: HashCode::rand(), 1361 amount: decimal("10"), 1362 exchange_base_url: url("https://exchange.test.com/"), 1363 metadata: None, 1364 wtid: ShortHashCode::rand(), 1365 creditor: magnet_payto("payto://iban/HU02162000031000164800000000?receiver-name=name"), 1366 }; 1367 let now = now_sql_stable_ts(); 1368 let later = now + Span::new().hours(2); 1369 // Insert 1370 assert_eq!( 1371 make_transfer(&pool, &req, &now).await.expect("transfer"), 1372 TransferResult::Success { 1373 id: 1, 1374 initiated_at: now 1375 } 1376 ); 1377 // Idempotent 1378 assert_eq!( 1379 make_transfer(&pool, &req, &later).await.expect("transfer"), 1380 TransferResult::Success { 1381 id: 1, 1382 initiated_at: now 1383 } 1384 ); 1385 // Request UID reuse 1386 assert_eq!( 1387 make_transfer( 1388 &pool, 1389 &db::Transfer { 1390 wtid: ShortHashCode::rand(), 1391 ..req.clone() 1392 }, 1393 &now 1394 ) 1395 .await 1396 .expect("transfer"), 1397 TransferResult::RequestUidReuse 1398 ); 1399 // wtid reuse 1400 assert_eq!( 1401 make_transfer( 1402 &pool, 1403 &db::Transfer { 1404 request_uid: HashCode::rand(), 1405 ..req.clone() 1406 }, 1407 &now 1408 ) 1409 .await 1410 .expect("transfer"), 1411 TransferResult::WtidReuse 1412 ); 1413 // Many 1414 assert_eq!( 1415 make_transfer( 1416 &pool, 1417 &db::Transfer { 1418 request_uid: HashCode::rand(), 1419 wtid: ShortHashCode::rand(), 1420 ..req 1421 }, 1422 &later 1423 ) 1424 .await 1425 .expect("transfer"), 1426 TransferResult::Success { 1427 id: 2, 1428 initiated_at: later 1429 } 1430 ); 1431 1432 // Get 1433 assert!(db::transfer_by_id(&pool, 1).await.unwrap().is_some()); 1434 assert!(db::transfer_by_id(&pool, 2).await.unwrap().is_some()); 1435 assert!(db::transfer_by_id(&pool, 3).await.unwrap().is_none()); 1436 assert_eq!( 1437 db::transfer_page(&pool, &None, &Page::default()) 1438 .await 1439 .unwrap() 1440 .len(), 1441 2 1442 ); 1443 } 1444 1445 #[db_test] 1446 async fn bounce(mut db: PgConnection) { 1447 let amount = amount("HUF:10"); 1448 let payto = magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"); 1449 let now = now_sql_stable_ts(); 1450 let date = Zoned::now().date(); 1451 1452 // Empty db 1453 assert!(db::pending_batch(&mut db, &now).await.unwrap().is_empty()); 1454 1455 // Insert 1456 assert_eq!( 1457 register_tx_in( 1458 &mut db, 1459 &TxIn { 1460 code: 13, 1461 amount, 1462 subject: "subject".into(), 1463 debtor: payto.clone(), 1464 value_date: date, 1465 status: TxStatus::Completed 1466 }, 1467 &None, 1468 &now 1469 ) 1470 .await 1471 .expect("register tx in"), 1472 AddIncomingResult::Success { 1473 new: true, 1474 pending: false, 1475 row_id: 1, 1476 valued_at: date 1477 } 1478 ); 1479 1480 // Bounce 1481 assert_eq!( 1482 register_bounce_tx_in( 1483 &mut db, 1484 &TxIn { 1485 code: 12, 1486 amount, 1487 subject: "subject".into(), 1488 debtor: payto.clone(), 1489 value_date: date, 1490 status: TxStatus::Completed 1491 }, 1492 "good reason", 1493 &now 1494 ) 1495 .await 1496 .expect("bounce"), 1497 BounceResult { 1498 tx_id: 2, 1499 tx_new: true, 1500 bounce_id: 1, 1501 bounce_new: true 1502 } 1503 ); 1504 // Idempotent 1505 assert_eq!( 1506 register_bounce_tx_in( 1507 &mut db, 1508 &TxIn { 1509 code: 12, 1510 amount, 1511 subject: "subject".into(), 1512 debtor: payto.clone(), 1513 value_date: date, 1514 status: TxStatus::Completed 1515 }, 1516 "good reason", 1517 &now 1518 ) 1519 .await 1520 .expect("bounce"), 1521 BounceResult { 1522 tx_id: 2, 1523 tx_new: false, 1524 bounce_id: 1, 1525 bounce_new: false 1526 } 1527 ); 1528 1529 // Bounce registered 1530 assert_eq!( 1531 register_bounce_tx_in( 1532 &mut db, 1533 &TxIn { 1534 code: 13, 1535 amount, 1536 subject: "subject".into(), 1537 debtor: payto.clone(), 1538 value_date: date, 1539 status: TxStatus::Completed 1540 }, 1541 "good reason", 1542 &now 1543 ) 1544 .await 1545 .expect("bounce"), 1546 BounceResult { 1547 tx_id: 1, 1548 tx_new: false, 1549 bounce_id: 2, 1550 bounce_new: true 1551 } 1552 ); 1553 // Idempotent registered 1554 assert_eq!( 1555 register_bounce_tx_in( 1556 &mut db, 1557 &TxIn { 1558 code: 13, 1559 amount, 1560 subject: "subject".into(), 1561 debtor: payto.clone(), 1562 value_date: date, 1563 status: TxStatus::Completed 1564 }, 1565 "good reason", 1566 &now 1567 ) 1568 .await 1569 .expect("bounce"), 1570 BounceResult { 1571 tx_id: 1, 1572 tx_new: false, 1573 bounce_id: 2, 1574 bounce_new: false 1575 } 1576 ); 1577 1578 // Batch 1579 assert_eq!( 1580 db::pending_batch(&mut db, &now).await.unwrap(), 1581 &[ 1582 Initiated { 1583 id: 1, 1584 amount, 1585 subject: "bounce: 12".into(), 1586 creditor: payto.clone() 1587 }, 1588 Initiated { 1589 id: 2, 1590 amount, 1591 subject: "bounce: 13".into(), 1592 creditor: payto 1593 } 1594 ] 1595 ); 1596 } 1597 1598 #[db_test] 1599 async fn status(mut db: PgConnection) { 1600 // Unknown transfer 1601 db::initiated_submit_permanent_failure(&mut db, 1, &Timestamp::now(), "msg") 1602 .await 1603 .unwrap(); 1604 db::initiated_submit_success(&mut db, 1, &Timestamp::now(), 12) 1605 .await 1606 .unwrap(); 1607 } 1608 1609 #[db_test] 1610 async fn batch(mut db: PgConnection, pool: PgPool) { 1611 let start = Timestamp::now(); 1612 let magnet_payto = 1613 magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"); 1614 1615 // Empty db 1616 let pendings = db::pending_batch(&mut db, &start) 1617 .await 1618 .expect("pending_batch"); 1619 assert_eq!(pendings.len(), 0); 1620 1621 // Some transfers 1622 for i in 0..3 { 1623 make_transfer( 1624 &pool, 1625 &db::Transfer { 1626 request_uid: HashCode::rand(), 1627 amount: decimal(format!("{}", i + 1)), 1628 exchange_base_url: url("https://exchange.test.com/"), 1629 metadata: None, 1630 wtid: ShortHashCode::rand(), 1631 creditor: magnet_payto.clone(), 1632 }, 1633 &Timestamp::now(), 1634 ) 1635 .await 1636 .expect("transfer"); 1637 } 1638 let pendings = db::pending_batch(&mut db, &start) 1639 .await 1640 .expect("pending_batch"); 1641 assert_eq!(pendings.len(), 3); 1642 1643 // Max 100 txs in batch 1644 for i in 0..100 { 1645 make_transfer( 1646 &pool, 1647 &db::Transfer { 1648 request_uid: HashCode::rand(), 1649 amount: decimal(format!("{}", i + 1)), 1650 exchange_base_url: url("https://exchange.test.com/"), 1651 metadata: None, 1652 wtid: ShortHashCode::rand(), 1653 creditor: magnet_payto.clone(), 1654 }, 1655 &Timestamp::now(), 1656 ) 1657 .await 1658 .expect("transfer"); 1659 } 1660 let pendings = db::pending_batch(&mut db, &start) 1661 .await 1662 .expect("pending_batch"); 1663 assert_eq!(pendings.len(), 100); 1664 1665 // Skip uploaded 1666 for i in 0..=10 { 1667 db::initiated_submit_success(&mut db, i, &Timestamp::now(), i) 1668 .await 1669 .expect("status success"); 1670 } 1671 let pendings = db::pending_batch(&mut db, &start) 1672 .await 1673 .expect("pending_batch"); 1674 assert_eq!(pendings.len(), 93); 1675 1676 // Skip failed 1677 for i in 0..=10 { 1678 db::initiated_submit_permanent_failure(&mut db, 10 + i, &Timestamp::now(), "failure") 1679 .await 1680 .expect("status failure"); 1681 } 1682 let pendings = db::pending_batch(&mut db, &start) 1683 .await 1684 .expect("pending_batch"); 1685 assert_eq!(pendings.len(), 83); 1686 } 1687 }