taler-rust

GNU Taler code in Rust. Largely core banking integrations.
Log | Files | Refs | Submodules | README | LICENSE

worker.rs (23323B)


      1 /*
      2   This file is part of TALER
      3   Copyright (C) 2025, 2026 Taler Systems SA
      4 
      5   TALER is free software; you can redistribute it and/or modify it under the
      6   terms of the GNU Affero General Public License as published by the Free Software
      7   Foundation; either version 3, or (at your option) any later version.
      8 
      9   TALER is distributed in the hope that it will be useful, but WITHOUT ANY
     10   WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR
     11   A PARTICULAR PURPOSE.  See the GNU Affero General Public License for more details.
     12 
     13   You should have received a copy of the GNU Affero General Public License along with
     14   TALER; see the file COPYING.  If not, see <http://www.gnu.org/licenses/>
     15 */
     16 
     17 use std::time::Duration;
     18 
     19 use failure_injection::{InjectedErr, fail_point};
     20 use http_client::ApiErr;
     21 use jiff::Timestamp;
     22 use sqlx::{Acquire as _, PgConnection, PgPool, postgres::PgListener};
     23 use taler_api::subject::{self, IncomingSubject, parse_incoming_unstructured};
     24 use taler_common::{
     25     ExpoBackoffDecorr,
     26     config::Config,
     27     types::amount::{self, Currency},
     28 };
     29 use tokio::sync::Notify;
     30 use tracing::{debug, error, info, trace, warn};
     31 
     32 use crate::{
     33     config::{AccountType, WorkerCfg},
     34     constants::SYNC_CURSOR_KEY,
     35     cyclos_api::{
     36         api::{CyclosAuth, CyclosErr},
     37         client::Client,
     38         types::{AccountKind, HistoryItem, InputError, NotFoundError, OrderBy},
     39     },
     40     db::{
     41         self, AddIncomingResult, ChargebackFailureResult, RegisterResult, TxIn, TxOut, TxOutKind,
     42         kv_get, kv_set,
     43     },
     44     notification::watch_notification,
     45 };
     46 
     47 #[derive(Debug, thiserror::Error)]
     48 pub enum WorkerError {
     49     #[error(transparent)]
     50     Db(#[from] sqlx::Error),
     51     #[error(transparent)]
     52     Api(#[from] ApiErr<CyclosErr>),
     53     #[error("Another worker is running concurrently")]
     54     Concurrency,
     55     #[error(transparent)]
     56     Injected(#[from] InjectedErr),
     57 }
     58 
     59 pub type WorkerResult = Result<(), WorkerError>;
     60 
     61 /// Retry an operation, not a terminated task. A closed pool cannot recover,
     62 /// and injected interruptions must remain visible to the caller/supervisor.
     63 fn retry_delay(
     64     err: &WorkerError,
     65     jitter: &mut ExpoBackoffDecorr,
     66     frequency: Duration,
     67 ) -> Option<Duration> {
     68     match err {
     69         WorkerError::Db(sqlx::Error::PoolClosed) | WorkerError::Injected(_) => None,
     70         WorkerError::Concurrency => Some(jitter.backoff().max(Duration::from_secs(15))),
     71         WorkerError::Api(e)
     72             if matches!(*e.err, CyclosErr::Input(InputError::Validation { .. })) =>
     73         {
     74             // Invalid input must not be resubmitted at notification speed.
     75             Some(jitter.backoff().max(frequency))
     76         }
     77         _ => Some(jitter.backoff()),
     78     }
     79 }
     80 
     81 pub async fn run_worker(
     82     cfg: &Config,
     83     pool: &PgPool,
     84     client: &http_client::Client,
     85     transient: bool,
     86 ) -> anyhow::Result<()> {
     87     let cfg = WorkerCfg::parse(cfg)?;
     88     let client = Client {
     89         client,
     90         api_url: &cfg.host.api_url,
     91         auth: &CyclosAuth::Basic {
     92             username: cfg.host.username,
     93             password: cfg.host.password,
     94         },
     95     };
     96     if transient {
     97         let mut conn = pool.acquire().await?;
     98         Worker {
     99             client: &client,
    100             db: &mut conn,
    101             account_type_id: *cfg.account_type_id,
    102             payment_type_id: *cfg.payment_type_id,
    103             account_type: cfg.account_type,
    104             currency: cfg.currency,
    105         }
    106         .run()
    107         .await?;
    108         return Ok(());
    109     }
    110 
    111     let notification = Notify::new();
    112 
    113     let watcher = async {
    114         watch_notification(&client, &notification).await;
    115     };
    116     let worker = async {
    117         let mut jitter = ExpoBackoffDecorr::default();
    118         loop {
    119             let res: WorkerResult = async {
    120                 let db = &mut PgListener::connect_with(pool).await?;
    121                 db.listen_all(["transfer"]).await?;
    122                 info!(target: "worker", "database listener connected");
    123                 loop {
    124                     let result = Worker {
    125                         client: &client,
    126                         db: db.acquire().await?,
    127                         account_type_id: *cfg.account_type_id,
    128                         payment_type_id: *cfg.payment_type_id,
    129                         account_type: cfg.account_type,
    130                         currency: cfg.currency,
    131                     }
    132                     .run()
    133                     .await;
    134                     match result {
    135                         Ok(()) => jitter.reset(),
    136                         Err(err @ WorkerError::Db(_)) => return Err(err),
    137                         Err(err) => {
    138                             let Some(delay) = retry_delay(&err, &mut jitter, cfg.frequency) else {
    139                                 return Err(err);
    140                             };
    141                             error!(target: "worker", ?delay, "synchronization pass failed: {err}");
    142                             tokio::time::sleep(delay).await;
    143                             continue;
    144                         }
    145                     }
    146                     tokio::select! {
    147                         _ = tokio::time::sleep(cfg.frequency) => {
    148                             info!(target: "worker", "running at frequency");
    149                         }
    150                         res = db.try_recv() => {
    151                             let mut ntf = res?;
    152                             while let Some(n) = ntf {
    153                                 debug!(target: "worker", "notification from {}", n.channel());
    154                                 ntf = db.next_buffered();
    155                             }
    156                         }
    157                         _ = notification.notified() => {
    158                             info!(target: "worker", "running at notification trigger");
    159                         }
    160                     }
    161                 }
    162             }
    163             .await;
    164             let err = res.unwrap_err();
    165             let Some(delay) = retry_delay(&err, &mut jitter, cfg.frequency) else {
    166                 break Err(err);
    167             };
    168             error!(target: "worker", ?delay, "database session failed: {err}");
    169             tokio::time::sleep(delay).await;
    170         }
    171     };
    172     // The optional notification listener reconnects independently. Polling
    173     // continues during its outages; a terminal worker result is not swallowed.
    174     tokio::select! {
    175         result = worker => result?,
    176         _ = watcher => unreachable!("notification listener does not return"),
    177     }
    178 }
    179 
    180 pub struct Worker<'a> {
    181     pub client: &'a Client<'a>,
    182     pub db: &'a mut PgConnection,
    183     pub currency: Currency,
    184     pub account_type_id: i64,
    185     pub payment_type_id: i64,
    186     pub account_type: AccountType,
    187 }
    188 
    189 impl Worker<'_> {
    190     /// Run a single worker pass
    191     pub async fn run(&mut self) -> WorkerResult {
    192         // Some worker operations are not idempotent, therefore it's not safe to have multiple worker
    193         // running concurrently. We use a global Postgres advisory lock to prevent it.
    194         if !db::worker_lock(self.db).await? {
    195             return Err(WorkerError::Concurrency);
    196         }
    197 
    198         // Sync transactions
    199         let mut cursor: Timestamp = kv_get(&mut *self.db, SYNC_CURSOR_KEY)
    200             .await?
    201             .unwrap_or_default();
    202 
    203         loop {
    204             let page = self
    205                 .client
    206                 .history(self.account_type_id, OrderBy::DateAsc, 0, Some(cursor))
    207                 .await?;
    208             for transfer in page.page {
    209                 if transfer.date > cursor {
    210                     cursor = transfer.date;
    211                 }
    212                 let tx = extract_tx_info(transfer);
    213                 match tx {
    214                     Tx::In(tx_in) => self.ingest_in(tx_in).await?,
    215                     Tx::Out(tx_out) => self.ingest_out(tx_out).await?,
    216                 }
    217             }
    218 
    219             kv_set(&mut *self.db, SYNC_CURSOR_KEY, &cursor).await?;
    220 
    221             if !page.has_next_page {
    222                 break;
    223             }
    224         }
    225 
    226         // Send transactions
    227         let start = Timestamp::now();
    228         loop {
    229             let batch = db::pending_batch(&mut *self.db, &start).await?;
    230             if batch.is_empty() {
    231                 break;
    232             }
    233             for initiated in batch {
    234                 debug!(target: "worker", "send tx {initiated}");
    235                 let res = self
    236                     .client
    237                     .direct_payment(
    238                         initiated.creditor_id,
    239                         self.payment_type_id,
    240                         initiated.amount,
    241                         &initiated.subject,
    242                     )
    243                     .await;
    244                 fail_point("direct-payment")?;
    245                 match res {
    246                     Ok(tx) => {
    247                         // Update transaction status, on failure the initiated transaction will be orphan
    248                         db::initiated_submit_success(
    249                             &mut *self.db,
    250                             initiated.id,
    251                             &tx.date,
    252                             tx.id.0,
    253                         )
    254                         .await?;
    255                         trace!(target: "worker", "init tx {}", tx.id);
    256                     }
    257                     Err(e) => {
    258                         let msg = match &*e.err {
    259                             CyclosErr::Unknown(NotFoundError { entity_type, key }) => {
    260                                 format!("unknown {entity_type} {key}")
    261                             }
    262                             CyclosErr::Forbidden(err) => err.to_string(),
    263                             _ => return Err(e.into()),
    264                         };
    265                         // TODO is permission should be considered are hard or soft failure ?
    266                         db::initiated_submit_permanent_failure(&mut *self.db, initiated.id, &msg)
    267                             .await?;
    268                         error!(target: "worker", "initiated failure {initiated}: {msg}");
    269                     }
    270                 }
    271             }
    272         }
    273         Ok(())
    274     }
    275 
    276     /// Ingest an incoming transaction
    277     async fn ingest_in(&mut self, tx: TxIn) -> WorkerResult {
    278         match self.account_type {
    279             AccountType::Exchange => {
    280                 let transfer = self.client.transfer(tx.transfer_id).await?;
    281                 let bounce = async |db: &mut PgConnection,
    282                                     reason: &str|
    283                        -> Result<(), WorkerError> {
    284                     // Fetch existing transaction
    285                     if let Some(chargeback) = transfer.charged_back_by {
    286                         let res = db::register_bounced_tx_in(
    287                             db,
    288                             &tx,
    289                             *chargeback.id,
    290                             reason,
    291                             &Timestamp::now(),
    292                         )
    293                         .await?;
    294                         if res.tx_new {
    295                             info!(target: "worker",
    296                                 "in {tx} bounced (recovered) in {}: {reason}", chargeback.id
    297                             );
    298                         } else {
    299                             trace!(target: "worker",
    300                                 "in {tx} already seen and bounced in {}: {reason}",chargeback.id
    301                             );
    302                         }
    303                     } else if !transfer.can_chargeback {
    304                         match db::register_tx_in(db, &tx, &None, &Timestamp::now()).await? {
    305                             AddIncomingResult::Success { new, .. } => {
    306                                 if new {
    307                                     warn!(target: "worker", "in {tx} cannot bounce: {reason}");
    308                                 } else {
    309                                     trace!(target: "worker", "in {tx} already seen and cannot bounce ");
    310                                 }
    311                             }
    312                             AddIncomingResult::ReservePubReuse
    313                             | AddIncomingResult::UnknownMapping
    314                             | AddIncomingResult::MappingReuse => unreachable!(),
    315                         }
    316                     } else {
    317                         let chargeback_id = self.client.chargeback(*transfer.id).await?;
    318                         fail_point("chargeback")?;
    319                         let res = db::register_bounced_tx_in(
    320                             db,
    321                             &tx,
    322                             chargeback_id,
    323                             reason,
    324                             &Timestamp::now(),
    325                         )
    326                         .await?;
    327                         if res.tx_new {
    328                             info!(target: "worker", "in {tx} bounced in {chargeback_id}: {reason}");
    329                         } else {
    330                             trace!(target: "worker", "in {tx} already seen and bounced in {chargeback_id}: {reason}");
    331                         }
    332                     }
    333                     Ok(())
    334                 };
    335                 if let Some(chargeback) = transfer.chargeback_of {
    336                     // This a chargeback of one of our transaction, if we bounce we might enter a loop
    337                     match db::initiated_chargeback_failure(&mut *self.db, *chargeback.id).await? {
    338                         ChargebackFailureResult::Unknown => {
    339                             trace!(target: "worker", "initiated failure unknown: charged back")
    340                         }
    341                         ChargebackFailureResult::Known(initiated) => {
    342                             error!(target: "worker", "initiated failure {initiated}: charged back")
    343                         }
    344                         ChargebackFailureResult::Idempotent(initiated) => {
    345                             trace!(target: "worker", "initiated failure {initiated} already seen: charged back")
    346                         }
    347                     }
    348                     // Sill register the incoming transaction as an incoming one
    349                     match db::register_tx_in(self.db, &tx, &None, &Timestamp::now()).await? {
    350                         AddIncomingResult::Success { new, .. } => {
    351                             if new {
    352                                 info!(target: "worker", "in {tx} chargeback");
    353                             } else {
    354                                 trace!(target: "worker", "in {tx} chargeback already seen");
    355                             }
    356                         }
    357                         AddIncomingResult::ReservePubReuse
    358                         | AddIncomingResult::UnknownMapping
    359                         | AddIncomingResult::MappingReuse => unreachable!(),
    360                     }
    361 
    362                     return Ok(());
    363                 }
    364                 match parse_incoming_unstructured(&tx.subject) {
    365                     Ok(subject) => {
    366                         match subject {
    367                             IncomingSubject::Key(subject) => {
    368                                 match db::register_tx_in(
    369                                     self.db,
    370                                     &tx,
    371                                     &Some(subject),
    372                                     &Timestamp::now(),
    373                                 )
    374                                 .await?
    375                                 {
    376                                     AddIncomingResult::Success { new, .. } => {
    377                                         if new {
    378                                             info!(target: "worker", "in {tx}");
    379                                         } else {
    380                                             trace!(target: "worker", "in {tx} already seen");
    381                                         }
    382                                     }
    383                                     AddIncomingResult::ReservePubReuse => {
    384                                         bounce(self.db, "reserve pub reuse").await?
    385                                     }
    386                                     AddIncomingResult::UnknownMapping => {
    387                                         bounce(self.db, "unknown mapping").await?
    388                                     }
    389                                     AddIncomingResult::MappingReuse => {
    390                                         bounce(self.db, "mapping reuse").await?
    391                                     }
    392                                 }
    393                             }
    394                             IncomingSubject::AdminBalanceAdjust => {
    395                                 // TODO bounce or skip ?
    396                             }
    397                         }
    398                     }
    399                     Err(e) => bounce(self.db, &e.to_string()).await?,
    400                 }
    401             }
    402             AccountType::Normal => {
    403                 match db::register_tx_in(self.db, &tx, &None, &Timestamp::now()).await? {
    404                     AddIncomingResult::Success { new, .. } => {
    405                         if new {
    406                             info!(target: "worker", "in {tx}");
    407                         } else {
    408                             trace!(target: "worker", "in {tx} already seen");
    409                         }
    410                     }
    411                     AddIncomingResult::ReservePubReuse
    412                     | AddIncomingResult::UnknownMapping
    413                     | AddIncomingResult::MappingReuse => unreachable!(),
    414                 }
    415             }
    416         }
    417         Ok(())
    418     }
    419 
    420     async fn ingest_out(&mut self, tx: TxOut) -> WorkerResult {
    421         match self.account_type {
    422             AccountType::Exchange => {
    423                 let transfer = self.client.transfer(tx.transfer_id).await?;
    424 
    425                 if transfer.charged_back_by.is_some() {
    426                     match db::initiated_chargeback_failure(&mut *self.db, *transfer.id).await? {
    427                         ChargebackFailureResult::Unknown => {
    428                             trace!(target: "worker", "initiated failure unknown: charged back")
    429                         }
    430                         ChargebackFailureResult::Known(initiated) => {
    431                             error!(target: "worker", "initiated failure {initiated}: charged back")
    432                         }
    433                         ChargebackFailureResult::Idempotent(initiated) => {
    434                             trace!(target: "worker", "initiated failure {initiated} already seen: charged back")
    435                         }
    436                     }
    437                 }
    438 
    439                 let kind = if let Ok(subject) = subject::parse_outgoing(&tx.subject) {
    440                     TxOutKind::Talerable(subject)
    441                 } else if let Some(chargeback) = &transfer.chargeback_of {
    442                     TxOutKind::Bounce(*chargeback.id)
    443                 } else {
    444                     TxOutKind::Simple
    445                 };
    446 
    447                 let res = db::register_tx_out(self.db, &tx, &kind, &Timestamp::now()).await?;
    448                 match res.result {
    449                     RegisterResult::idempotent => match kind {
    450                         TxOutKind::Simple => {
    451                             trace!(target: "worker", "out malformed {tx} already seen")
    452                         }
    453                         TxOutKind::Bounce(_) => {
    454                             trace!(target: "worker", "out bounce {tx} already seen")
    455                         }
    456                         TxOutKind::Talerable(_) => {
    457                             trace!(target: "worker", "out {tx} already seen")
    458                         }
    459                     },
    460                     RegisterResult::known => match kind {
    461                         TxOutKind::Simple => {
    462                             warn!(target: "worker", "out malformed {tx}")
    463                         }
    464                         TxOutKind::Bounce(_) => {
    465                             info!(target: "worker", "out bounce {tx}")
    466                         }
    467                         TxOutKind::Talerable(_) => {
    468                             info!(target: "worker", "out {tx}")
    469                         }
    470                     },
    471                     RegisterResult::recovered => match kind {
    472                         TxOutKind::Simple => {
    473                             warn!(target: "worker", "out malformed (recovered) {tx}")
    474                         }
    475                         TxOutKind::Bounce(_) => {
    476                             warn!(target: "worker", "out bounce (recovered) {tx}")
    477                         }
    478                         TxOutKind::Talerable(_) => {
    479                             warn!(target: "worker", "out (recovered) {tx}")
    480                         }
    481                     },
    482                 }
    483             }
    484             AccountType::Normal => {
    485                 let res = db::register_tx_out(self.db, &tx, &TxOutKind::Simple, &Timestamp::now())
    486                     .await?;
    487                 match res.result {
    488                     RegisterResult::idempotent => {
    489                         trace!(target: "worker", "out {tx} already seen");
    490                     }
    491                     RegisterResult::known => {
    492                         info!(target: "worker", "out {tx}");
    493                     }
    494                     RegisterResult::recovered => {
    495                         warn!(target: "worker", "out (recovered) {tx}");
    496                     }
    497                 }
    498             }
    499         }
    500         Ok(())
    501     }
    502 }
    503 
    504 pub enum Tx {
    505     In(TxIn),
    506     Out(TxOut),
    507 }
    508 
    509 pub fn extract_tx_info(tx: HistoryItem) -> Tx {
    510     let amount = amount::decimal(tx.amount.trim_start_matches('-'));
    511     let (id, name) = match tx.related_account.kind {
    512         AccountKind::System => (tx.related_account.ty.id, tx.related_account.ty.name),
    513         AccountKind::User { user } => (user.id, user.display),
    514     };
    515     if tx.amount.starts_with('-') {
    516         Tx::Out(TxOut {
    517             transfer_id: *tx.id,
    518             tx_id: tx.transaction.map(|it| *it.id),
    519             amount,
    520             subject: tx.description.unwrap_or_default(),
    521             creditor_id: *id,
    522             creditor_name: name,
    523             valued_at: tx.date,
    524         })
    525     } else {
    526         Tx::In(TxIn {
    527             transfer_id: *tx.id,
    528             tx_id: tx.transaction.map(|it| *it.id),
    529             amount,
    530             subject: tx.description.unwrap_or_default(),
    531             debtor_id: *id,
    532             debtor_name: name,
    533             valued_at: tx.date,
    534         })
    535     }
    536 }
    537 
    538 #[cfg(test)]
    539 mod retry_tests {
    540     use super::*;
    541 
    542     #[test]
    543     fn database_outages_retry_but_terminal_interruptions_propagate() {
    544         let mut jitter = ExpoBackoffDecorr::default();
    545         let err = WorkerError::Db(sqlx::Error::Io(
    546             std::io::ErrorKind::ConnectionRefused.into(),
    547         ));
    548         let delay = retry_delay(&err, &mut jitter, Duration::from_secs(60)).unwrap();
    549         assert!(delay >= Duration::from_millis(400));
    550         assert!(delay <= Duration::from_secs(30));
    551         assert!(
    552             retry_delay(
    553                 &WorkerError::Db(sqlx::Error::PoolClosed),
    554                 &mut jitter,
    555                 Duration::ZERO
    556             )
    557             .is_none()
    558         );
    559         assert!(
    560             retry_delay(
    561                 &WorkerError::Injected(InjectedErr("worker interrupted")),
    562                 &mut jitter,
    563                 Duration::ZERO
    564             )
    565             .is_none()
    566         );
    567         assert!(
    568             retry_delay(&WorkerError::Concurrency, &mut jitter, Duration::ZERO).unwrap()
    569                 >= Duration::from_secs(15)
    570         );
    571     }
    572 }