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 }