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