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, ¬ification).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 }