depolymerization

wire gateway for Bitcoin/Ethereum
Log | Files | Refs | Submodules | README | LICENSE

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 }