taler-rust

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

worker.rs (29841B)


      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::{num::ParseIntError, time::Duration};
     18 
     19 use aws_lc_rs::signature::EcdsaKeyPair;
     20 use failure_injection::{InjectedErr, fail_point};
     21 use http_client::ApiErr;
     22 use jiff::{Timestamp, Zoned, civil::Date};
     23 use sqlx::{Acquire as _, PgConnection, PgPool, postgres::PgListener};
     24 use taler_api::subject::{self, IncomingSubject, parse_incoming_unstructured};
     25 use taler_common::{
     26     ExpoBackoffDecorr,
     27     config::Config,
     28     types::{
     29         amount::{self},
     30         iban::IBAN,
     31     },
     32 };
     33 use tracing::{debug, error, info, trace, warn};
     34 
     35 use crate::{
     36     FullHuPayto, HuIban,
     37     config::{AccountType, WorkerCfg},
     38     db::{self, AddIncomingResult, Initiated, RegisterResult, TxIn, TxOut, TxOutKind},
     39     magnet_api::{
     40         api::MagnetErr,
     41         client::{ApiClient, AuthClient},
     42         types::{Direction, Next, Order, TxDto, TxStatus},
     43     },
     44     setup,
     45 };
     46 
     47 // const TXS_CURSOR_KEY: &str = "txs_cursor"; TODO cursor is broken
     48 
     49 #[derive(Debug, thiserror::Error)]
     50 pub enum WorkerError {
     51     #[error(transparent)]
     52     Db(#[from] sqlx::Error),
     53     #[error(transparent)]
     54     Api(#[from] ApiErr<MagnetErr>),
     55     #[error("Another worker is running concurrently")]
     56     Concurrency,
     57     #[error(transparent)]
     58     Injected(#[from] InjectedErr),
     59 }
     60 
     61 pub type WorkerResult = Result<(), WorkerError>;
     62 
     63 /// Retry an operation, not a terminated task. A closed pool cannot recover,
     64 /// and injected interruptions must remain visible to the caller/supervisor.
     65 fn retry_delay(err: &WorkerError, jitter: &mut ExpoBackoffDecorr) -> Option<Duration> {
     66     match err {
     67         WorkerError::Db(sqlx::Error::PoolClosed) | WorkerError::Injected(_) => None,
     68         WorkerError::Concurrency => Some(jitter.backoff().max(Duration::from_secs(15))),
     69         _ => Some(jitter.backoff()),
     70     }
     71 }
     72 
     73 pub async fn run_worker(
     74     cfg: &Config,
     75     pool: &PgPool,
     76     client: &http_client::Client,
     77     transient: bool,
     78 ) -> anyhow::Result<()> {
     79     let cfg = WorkerCfg::parse(cfg)?;
     80     let keys = setup::load(&cfg)?;
     81     let client = AuthClient::new(client, &cfg.api_url, &cfg.consumer).upgrade(&keys.access_token);
     82 
     83     if transient {
     84         let mut conn = pool.acquire().await?;
     85         let account = client.account(cfg.payto.bban()).await?;
     86         Worker {
     87             client: &client,
     88             db: &mut conn,
     89             account_number: &account.number,
     90             account_code: account.code,
     91             key: &keys.signing_key,
     92             account_type: cfg.account_type,
     93             ignore_tx_before: cfg.ignore_tx_before,
     94             ignore_bounces_before: cfg.ignore_bounces_before,
     95         }
     96         .run()
     97         .await?;
     98         return Ok(());
     99     }
    100 
    101     let mut jitter = ExpoBackoffDecorr::default();
    102     // Account discovery is an API operation; retry it without reloading keys
    103     // or rebuilding the HTTP client.
    104     let account = loop {
    105         match client.account(cfg.payto.bban()).await {
    106             Ok(account) => break account,
    107             Err(err) => {
    108                 error!(target: "worker", "account lookup failed: {err}");
    109                 tokio::time::sleep(jitter.backoff()).await;
    110             }
    111         }
    112     };
    113     jitter.reset();
    114 
    115     loop {
    116         // Only database/session failures reconstruct the listener. API failures
    117         // retry a pass with the existing client and database session.
    118         let res: WorkerResult = async {
    119             let db = &mut PgListener::connect_with(pool).await?;
    120             db.listen_all(["transfer"]).await?;
    121             info!(target: "worker", "database listener connected");
    122             loop {
    123                 let result = Worker {
    124                     client: &client,
    125                     db: db.acquire().await?,
    126                     account_number: &account.number,
    127                     account_code: account.code,
    128                     key: &keys.signing_key,
    129                     account_type: cfg.account_type,
    130                     ignore_tx_before: cfg.ignore_tx_before,
    131                     ignore_bounces_before: cfg.ignore_bounces_before,
    132                 }
    133                 .run()
    134                 .await;
    135                 match result {
    136                     Ok(()) => jitter.reset(),
    137                     Err(err @ WorkerError::Db(_)) => return Err(err),
    138                     Err(err) => {
    139                         let Some(delay) = retry_delay(&err, &mut jitter) else {
    140                             return Err(err);
    141                         };
    142                         error!(target: "worker", ?delay, "synchronization pass failed: {err}");
    143                         tokio::time::sleep(delay).await;
    144                         continue;
    145                     }
    146                 }
    147                 // Notifications may accelerate successful polling, but never
    148                 // bypass the backoff after a failed operation.
    149                 match tokio::time::timeout(cfg.frequency, db.try_recv()).await {
    150                     Ok(res) => {
    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                     Err(_) => info!(target: "worker", "running at frequency"),
    158                 }
    159             }
    160         }
    161         .await;
    162         let err = res.unwrap_err();
    163         let Some(delay) = retry_delay(&err, &mut jitter) else {
    164             return Err(err.into());
    165         };
    166         error!(target: "worker", ?delay, "database session failed: {err}");
    167         tokio::time::sleep(delay).await;
    168     }
    169 }
    170 
    171 pub struct Worker<'a> {
    172     pub client: &'a ApiClient<'a>,
    173     pub db: &'a mut PgConnection,
    174     pub account_number: &'a str,
    175     pub account_code: u64,
    176     pub key: &'a EcdsaKeyPair,
    177     pub account_type: AccountType,
    178     pub ignore_tx_before: Option<Date>,
    179     pub ignore_bounces_before: Option<Date>,
    180 }
    181 
    182 impl Worker<'_> {
    183     /// Run a single worker pass
    184     pub async fn run(&mut self) -> WorkerResult {
    185         // Some worker operations are not idempotent, therefore it's not safe to have multiple worker
    186         // running concurrently. We use a global Postgres advisory lock to prevent it.
    187         if !db::worker_lock(self.db).await? {
    188             return Err(WorkerError::Concurrency);
    189         };
    190 
    191         // Sync transactions
    192         let mut next: Option<Next> = None; //kv_get(&mut *self.db, TXS_CURSOR_KEY).await?; TODO cursor logic is broken and cannot be stored & reused
    193         let mut all_final = true;
    194         let mut first = true;
    195         loop {
    196             let page = self
    197                 .client
    198                 .page_tx(
    199                     Direction::Both,
    200                     Order::Ascending,
    201                     100,
    202                     self.account_number,
    203                     &next,
    204                     first,
    205                 )
    206                 .await?;
    207             first = false;
    208             next = page.next;
    209             for item in page.list {
    210                 all_final &= item.tx.status.is_final();
    211                 let tx = extract_tx_info(item.tx);
    212                 match tx {
    213                     Tx::In(tx_in) => {
    214                         // We only register final successful incoming transactions
    215                         if tx_in.status != TxStatus::Completed {
    216                             debug!(target: "worker", "pending or failed in {tx_in}");
    217                             continue;
    218                         }
    219 
    220                         if let Some(before) = self.ignore_tx_before
    221                             && tx_in.value_date < before
    222                         {
    223                             debug!(target: "worker", "ignore in {tx_in}");
    224                             continue;
    225                         }
    226                         let bounce = async |db: &mut PgConnection,
    227                                             reason: &str|
    228                                -> Result<(), WorkerError> {
    229                             if let Some(before) = self.ignore_bounces_before
    230                                 && tx_in.value_date < before
    231                             {
    232                                 match db::register_tx_in(db, &tx_in, &None, &Timestamp::now())
    233                                     .await?
    234                                 {
    235                                     AddIncomingResult::Success { new, .. } => {
    236                                         if new {
    237                                             info!(target: "worker", "in  {tx_in} skip bounce: {reason}");
    238                                         } else {
    239                                             trace!(target: "worker", "in  {tx_in} already skip bounce ");
    240                                         }
    241                                     }
    242                                     AddIncomingResult::ReservePubReuse
    243                                     | AddIncomingResult::UnknownMapping
    244                                     | AddIncomingResult::MappingReuse => unreachable!(),
    245                                 }
    246                             } else {
    247                                 let res = db::register_bounce_tx_in(
    248                                     db,
    249                                     &tx_in,
    250                                     reason,
    251                                     &Timestamp::now(),
    252                                 )
    253                                 .await?;
    254 
    255                                 if res.tx_new {
    256                                     info!(target: "worker",
    257                                         "in  {tx_in} bounced in {}: {reason}",
    258                                         res.bounce_id
    259                                     );
    260                                 } else {
    261                                     trace!(target: "worker",
    262                                         "in  {tx_in} already seen and bounced in {}: {reason}",
    263                                         res.bounce_id
    264                                     );
    265                                 }
    266                             }
    267                             Ok(())
    268                         };
    269                         match self.account_type {
    270                             AccountType::Exchange => {
    271                                 match parse_incoming_unstructured(&tx_in.subject) {
    272                                     Ok(subject) => match subject {
    273                                         IncomingSubject::Key(subject) => {
    274                                             match db::register_tx_in(
    275                                                 self.db,
    276                                                 &tx_in,
    277                                                 &Some(subject),
    278                                                 &Timestamp::now(),
    279                                             )
    280                                             .await?
    281                                             {
    282                                                 AddIncomingResult::Success { new, .. } => {
    283                                                     if new {
    284                                                         info!(target: "worker", "in  {tx_in}");
    285                                                     } else {
    286                                                         trace!(target: "worker", "in  {tx_in} already seen");
    287                                                     }
    288                                                 }
    289                                                 AddIncomingResult::ReservePubReuse => {
    290                                                     bounce(self.db, "reserve pub reuse").await?
    291                                                 }
    292                                                 AddIncomingResult::UnknownMapping => {
    293                                                     bounce(self.db, "unknown mapping").await?
    294                                                 }
    295                                                 AddIncomingResult::MappingReuse => {
    296                                                     bounce(self.db, "mapping reuse").await?
    297                                                 }
    298                                             }
    299                                         }
    300                                         IncomingSubject::AdminBalanceAdjust => {
    301                                             // TODO bounce or skip ?
    302                                         }
    303                                     },
    304                                     Err(e) => bounce(self.db, &e.to_string()).await?,
    305                                 }
    306                             }
    307                             AccountType::Normal => {
    308                                 match db::register_tx_in(self.db, &tx_in, &None, &Timestamp::now())
    309                                     .await?
    310                                 {
    311                                     AddIncomingResult::Success { new, .. } => {
    312                                         if new {
    313                                             info!(target: "worker", "in  {tx_in}");
    314                                         } else {
    315                                             trace!(target: "worker", "in  {tx_in} already seen");
    316                                         }
    317                                     }
    318                                     AddIncomingResult::ReservePubReuse
    319                                     | AddIncomingResult::UnknownMapping
    320                                     | AddIncomingResult::MappingReuse => unreachable!(),
    321                                 }
    322                             }
    323                         }
    324                     }
    325                     Tx::Out(tx_out) => {
    326                         match tx_out.status {
    327                             TxStatus::ToBeRecorded => {
    328                                 self.recover_tx(&tx_out).await?;
    329                                 continue;
    330                             }
    331                             TxStatus::PendingFirstSignature
    332                             | TxStatus::PendingSecondSignature
    333                             | TxStatus::PendingProcessing
    334                             | TxStatus::Verified
    335                             | TxStatus::PartiallyCompleted
    336                             | TxStatus::UnderReview => {
    337                                 // Still pending
    338                                 debug!(target: "worker", "pending out {tx_out}");
    339                                 continue;
    340                             }
    341                             TxStatus::Rejected | TxStatus::Canceled | TxStatus::Completed => {}
    342                         }
    343                         match self.account_type {
    344                             AccountType::Exchange => {
    345                                 let kind = if let Ok(subject) =
    346                                     subject::parse_outgoing(&tx_out.subject)
    347                                 {
    348                                     TxOutKind::Talerable(subject)
    349                                 } else if let Ok(bounced) = parse_bounce_outgoing(&tx_out.subject) {
    350                                     TxOutKind::Bounce(bounced)
    351                                 } else {
    352                                     TxOutKind::Simple
    353                                 };
    354                                 if tx_out.status == TxStatus::Completed {
    355                                     let res = db::register_tx_out(
    356                                         self.db,
    357                                         &tx_out,
    358                                         &kind,
    359                                         &Timestamp::now(),
    360                                     )
    361                                     .await?;
    362                                     match res.result {
    363                                         RegisterResult::idempotent => match kind {
    364                                             TxOutKind::Simple => {
    365                                                 trace!(target: "worker", "out malformed {tx_out} already seen")
    366                                             }
    367                                             TxOutKind::Bounce(_) => {
    368                                                 trace!(target: "worker", "out bounce {tx_out} already seen")
    369                                             }
    370                                             TxOutKind::Talerable(_) => {
    371                                                 trace!(target: "worker", "out {tx_out} already seen")
    372                                             }
    373                                         },
    374                                         RegisterResult::known => match kind {
    375                                             TxOutKind::Simple => {
    376                                                 warn!(target: "worker", "out malformed {tx_out}")
    377                                             }
    378                                             TxOutKind::Bounce(_) => {
    379                                                 info!(target: "worker", "out bounce {tx_out}")
    380                                             }
    381                                             TxOutKind::Talerable(_) => {
    382                                                 info!(target: "worker", "out {tx_out}")
    383                                             }
    384                                         },
    385                                         RegisterResult::recovered => match kind {
    386                                             TxOutKind::Simple => {
    387                                                 warn!(target: "worker", "out malformed (recovered) {tx_out}")
    388                                             }
    389                                             TxOutKind::Bounce(_) => {
    390                                                 warn!(target: "worker", "out bounce (recovered) {tx_out}")
    391                                             }
    392                                             TxOutKind::Talerable(_) => {
    393                                                 warn!(target: "worker", "out (recovered) {tx_out}")
    394                                             }
    395                                         },
    396                                     }
    397                                 } else {
    398                                     let bounced = match kind {
    399                                         TxOutKind::Simple => None,
    400                                         TxOutKind::Bounce(bounced) => Some(bounced),
    401                                         TxOutKind::Talerable(_) => None,
    402                                     };
    403                                     let res = db::register_tx_out_failure(
    404                                         self.db,
    405                                         tx_out.code,
    406                                         bounced,
    407                                         &Timestamp::now(),
    408                                     )
    409                                     .await?;
    410                                     if let Some(id) = res.initiated_id {
    411                                         if res.new {
    412                                             error!(target: "worker", "out failure {id} {tx_out}");
    413                                         } else {
    414                                             trace!(target: "worker", "out failure {id} {tx_out} already seen");
    415                                         }
    416                                     }
    417                                 }
    418                             }
    419                             AccountType::Normal => {
    420                                 if tx_out.status == TxStatus::Completed {
    421                                     let res = db::register_tx_out(
    422                                         self.db,
    423                                         &tx_out,
    424                                         &TxOutKind::Simple,
    425                                         &Timestamp::now(),
    426                                     )
    427                                     .await?;
    428                                     match res.result {
    429                                         RegisterResult::idempotent => {
    430                                             trace!(target: "worker", "out {tx_out} already seen");
    431                                         }
    432                                         RegisterResult::known => {
    433                                             info!(target: "worker", "out {tx_out}");
    434                                         }
    435                                         RegisterResult::recovered => {
    436                                             warn!(target: "worker", "out (recovered) {tx_out}");
    437                                         }
    438                                     }
    439                                 } else {
    440                                     let res = db::register_tx_out_failure(
    441                                         self.db,
    442                                         tx_out.code,
    443                                         None,
    444                                         &Timestamp::now(),
    445                                     )
    446                                     .await?;
    447                                     if let Some(id) = res.initiated_id {
    448                                         if res.new {
    449                                             error!(target: "worker", "out failure {id} {tx_out}");
    450                                         } else {
    451                                             trace!(target: "worker", "out failure {id} {tx_out} already seen");
    452                                         }
    453                                     }
    454                                 }
    455                             }
    456                         }
    457                     }
    458                 }
    459             }
    460 
    461             if let Some(_next) = &next {
    462                 // Update in db cursor only if all previous transactions where final
    463                 if all_final {
    464                     // debug!(target: "worker", "advance cursor {next:?}");
    465                     // kv_set(&mut *self.db, TXS_CURSOR_KEY, &next).await?; TODO cursor is broken
    466                 }
    467             } else {
    468                 break;
    469             }
    470         }
    471 
    472         // Send transactions
    473         let start = Timestamp::now();
    474         let now = Zoned::now();
    475         loop {
    476             let batch = db::pending_batch(&mut *self.db, &start).await?;
    477             if batch.is_empty() {
    478                 break;
    479             }
    480             for tx in batch {
    481                 debug!(target: "worker", "send tx {tx}");
    482                 self.init_tx(&tx, &now).await?;
    483             }
    484         }
    485         Ok(())
    486     }
    487 
    488     /// Try to sign an unsigned initiated transaction
    489     pub async fn recover_tx(&mut self, tx: &TxOut) -> WorkerResult {
    490         if db::initiated_exists_for_code(&mut *self.db, tx.code)
    491             .await?
    492             .is_some()
    493         {
    494             // Known initiated we submit it
    495             assert_eq!(tx.amount.frac, 0);
    496             self.submit_tx(
    497                 tx.code,
    498                 -(tx.amount.val as f64),
    499                 &tx.value_date,
    500                 tx.creditor.bban(),
    501             )
    502             .await?;
    503         } else {
    504             // The transaction is unknown (we failed after creating it and before storing it in the db)
    505             // we delete it
    506             self.client.delete_tx(tx.code).await?;
    507             debug!(target: "worker", "out {}: delete uncompleted orphan", tx.code);
    508         }
    509 
    510         Ok(())
    511     }
    512 
    513     /// Create and sign a forint transfer
    514     pub async fn init_tx(&mut self, tx: &Initiated, now: &Zoned) -> WorkerResult {
    515         trace!(target: "worker", "init tx {tx}");
    516         assert_eq!(tx.amount.frac, 0);
    517         let date = now.date();
    518         // Initialize the new transaction, on failure an orphan initiated transaction can be created
    519         let res = self
    520             .client
    521             .init_tx(
    522                 self.account_code,
    523                 tx.amount.val as f64,
    524                 &tx.subject,
    525                 &date,
    526                 &tx.creditor.name,
    527                 tx.creditor.bban(),
    528             )
    529             .await;
    530         fail_point("init-tx")?;
    531         let info = match res {
    532             // Check if succeeded
    533             Ok(info) => {
    534                 // Update transaction status, on failure the initiated transaction will be orphan
    535                 db::initiated_submit_success(&mut *self.db, tx.id, &Timestamp::now(), info.code)
    536                     .await?;
    537                 info
    538             }
    539             Err(e) => {
    540                 if let MagnetErr::Magnet(e) = &*e.err {
    541                     // Check if error is permanent
    542                     if matches!(
    543                         (e.error_code, e.short_message.as_str()),
    544                         (404, "BSZLA_NEM_TALALHATO") // Unknown account
    545                          | (409, "FORRAS_SZAMLA_ESZAMLA_EGYEZIK") // Same account
    546                     ) {
    547                         db::initiated_submit_permanent_failure(
    548                             &mut *self.db,
    549                             tx.id,
    550                             &Timestamp::now(),
    551                             &e.to_string(),
    552                         )
    553                         .await?;
    554                         error!(target: "worker", "initiated failure {tx}: {e}");
    555                         return WorkerResult::Ok(());
    556                     }
    557                 }
    558                 return Err(e.into());
    559             }
    560         };
    561         trace!(target: "worker", "init tx {}", info.code);
    562 
    563         // Sign transaction
    564         self.submit_tx(info.code, info.amount, &date, tx.creditor.bban())
    565             .await?;
    566         Ok(())
    567     }
    568 
    569     /** Submit an initiated forint transfer */
    570     pub async fn submit_tx(
    571         &mut self,
    572         tx_code: u64,
    573         amount: f64,
    574         date: &Date,
    575         creditor: &str,
    576     ) -> WorkerResult {
    577         debug!(target: "worker", "submit tx {tx_code}");
    578         fail_point("submit-tx")?;
    579         // Submit an initiated transaction, on failure we will retry
    580         match self
    581             .client
    582             .submit_tx(
    583                 self.key,
    584                 self.account_number,
    585                 tx_code,
    586                 amount,
    587                 date,
    588                 creditor,
    589             )
    590             .await
    591         {
    592             Ok(_) => Ok(()),
    593             Err(e) => {
    594                 if let MagnetErr::Magnet(e) = &*e.err {
    595                     // Check if soft failure
    596                     if matches!(
    597                         (e.error_code, e.short_message.as_str()),
    598                         (409, "TRANZAKCIO_ROSSZ_STATUS") // Already summited or cannot be signed
    599                     ) {
    600                         warn!(target: "worker", "submit tx {tx_code}: {e}");
    601                         return Ok(());
    602                     }
    603                 }
    604                 Err(e.into())
    605             }
    606         }
    607     }
    608 }
    609 
    610 pub enum Tx {
    611     In(TxIn),
    612     Out(TxOut),
    613 }
    614 
    615 pub fn extract_tx_info(tx: TxDto) -> Tx {
    616     // TODO amount from f64 without allocations
    617     let amount = amount::amount(format!("{}:{}", tx.currency, tx.amount.abs()));
    618     // TODO we should support non hungarian account and error handling
    619     let iban = if tx.counter_account.starts_with("HU") {
    620         let iban: IBAN = tx.counter_account.parse().unwrap();
    621         HuIban::try_from(iban).unwrap()
    622     } else {
    623         HuIban::from_bban(&tx.counter_account).unwrap()
    624     };
    625     let counter_account = FullHuPayto::new(iban, &tx.counter_name);
    626     if tx.amount.is_sign_positive() {
    627         Tx::In(TxIn {
    628             code: tx.code,
    629             amount,
    630             subject: tx.subject.unwrap_or_default(),
    631             debtor: counter_account,
    632             value_date: tx.value_date,
    633             status: tx.status,
    634         })
    635     } else {
    636         Tx::Out(TxOut {
    637             code: tx.code,
    638             amount,
    639             subject: tx.subject.unwrap_or_default(),
    640             creditor: counter_account,
    641             value_date: tx.value_date,
    642             status: tx.status,
    643         })
    644     }
    645 }
    646 
    647 #[derive(Debug, thiserror::Error)]
    648 pub enum BounceSubjectErr {
    649     #[error("missing parts")]
    650     MissingParts,
    651     #[error("not a bounce")]
    652     NotBounce,
    653     #[error("malformed bounced id: {0}")]
    654     Id(#[from] ParseIntError),
    655 }
    656 
    657 pub fn parse_bounce_outgoing(subject: &str) -> Result<u32, BounceSubjectErr> {
    658     let (prefix, id) = subject
    659         .rsplit_once(" ")
    660         .ok_or(BounceSubjectErr::MissingParts)?;
    661     if !prefix.starts_with("bounce") {
    662         return Err(BounceSubjectErr::NotBounce);
    663     }
    664     let id: u32 = id.parse()?;
    665     Ok(id)
    666 }
    667 
    668 #[cfg(test)]
    669 mod retry_tests {
    670     use super::*;
    671 
    672     #[test]
    673     fn database_outages_retry_but_terminal_interruptions_propagate() {
    674         let mut jitter = ExpoBackoffDecorr::default();
    675         let err = WorkerError::Db(sqlx::Error::Io(
    676             std::io::ErrorKind::ConnectionRefused.into(),
    677         ));
    678         let delay = retry_delay(&err, &mut jitter).unwrap();
    679         assert!(delay >= Duration::from_millis(400));
    680         assert!(delay <= Duration::from_secs(30));
    681         assert!(retry_delay(&WorkerError::Db(sqlx::Error::PoolClosed), &mut jitter).is_none());
    682         assert!(
    683             retry_delay(
    684                 &WorkerError::Injected(InjectedErr("worker interrupted")),
    685                 &mut jitter
    686             )
    687             .is_none()
    688         );
    689         assert!(
    690             retry_delay(&WorkerError::Concurrency, &mut jitter).unwrap() >= Duration::from_secs(15)
    691         );
    692     }
    693 }