api.rs (19566B)
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 jiff::Timestamp; 18 use prometheus_client::{metrics::gauge::Gauge, registry::Registry}; 19 use sqlx::PgPool; 20 use taler_api::{ 21 api::{ 22 TalerApi, observability::Observability, prepared::PreparedTransfer, revenue::Revenue, 23 wire::WireGateway, 24 }, 25 error::{ApiResult, failure_code}, 26 subject::{IncomingKey, fmt_in_subject}, 27 }; 28 use taler_common::{ 29 api::{ 30 params::{History, Page}, 31 prepared::{ 32 RegistrationRequest, RegistrationResponse, SubjectFormat, TransferSubject, 33 Unregistration, 34 }, 35 revenue::RevenueIncomingHistory, 36 wire::{ 37 AddIncomingRequest, AddIncomingResponse, AddKycauthRequest, AddMappedRequest, 38 IncomingHistory, OutgoingHistory, TransferList, TransferRequest, TransferResponse, 39 TransferState, TransferStatus, 40 }, 41 }, 42 db::IncomingType, 43 error_code::ErrorCode, 44 types::{amount::Currency, time::TalerTimestamp, utils::date_to_utc_ts}, 45 }; 46 use tokio::sync::watch::Sender; 47 48 use crate::{ 49 FullHuPayto, 50 constants::CURR, 51 db::{self, AddIncomingResult, Transfer, TxInAdmin}, 52 }; 53 54 pub struct MagnetApi { 55 pub pool: sqlx::PgPool, 56 pub payto: FullHuPayto, 57 pub in_channel: Sender<i64>, 58 pub taler_in_channel: Sender<i64>, 59 pub out_channel: Sender<i64>, 60 pub taler_out_channel: Sender<i64>, 61 metrics: Metrics, 62 registry: Registry, 63 } 64 65 #[derive(Default)] 66 struct Metrics { 67 db_access: Gauge, 68 } 69 70 impl Metrics { 71 pub fn registry(&self) -> Registry { 72 let mut registry = Registry::default(); 73 74 registry.register( 75 "db_access", 76 "Whether the last database metrics refresh succeeded", 77 self.db_access.clone(), 78 ); 79 registry 80 } 81 82 pub async fn sync(&self, db: &PgPool) { 83 let test = sqlx::query("SELECT 1").fetch_one(db).await.is_ok(); 84 self.db_access.set(if test { 1 } else { 0 }); 85 } 86 } 87 88 impl MagnetApi { 89 pub async fn start(pool: sqlx::PgPool, payto: FullHuPayto) -> Self { 90 let in_channel = Sender::new(0); 91 let taler_in_channel = Sender::new(0); 92 let out_channel = Sender::new(0); 93 let taler_out_channel = Sender::new(0); 94 95 let metrics = Metrics::default(); 96 97 let tmp = Self { 98 pool: pool.clone(), 99 payto, 100 in_channel: in_channel.clone(), 101 taler_in_channel: taler_in_channel.clone(), 102 out_channel: out_channel.clone(), 103 taler_out_channel: taler_out_channel.clone(), 104 registry: metrics.registry(), 105 metrics, 106 }; 107 tokio::spawn(db::notification_listener( 108 pool, 109 in_channel, 110 taler_in_channel, 111 out_channel, 112 taler_out_channel, 113 )); 114 tmp 115 } 116 } 117 118 impl TalerApi for MagnetApi { 119 fn currency(&self) -> Currency { 120 CURR 121 } 122 123 fn implementation(&self) -> &'static str { 124 "urn:net:taler:specs:taler-magnet-bank:taler-rust" 125 } 126 } 127 128 impl WireGateway for MagnetApi { 129 async fn transfer(&self, req: TransferRequest) -> ApiResult<TransferResponse> { 130 let creditor = FullHuPayto::try_from(&req.credit_account)?; 131 let result = db::make_transfer( 132 &self.pool, 133 &Transfer { 134 request_uid: req.request_uid, 135 wtid: req.wtid, 136 amount: req.amount.decimal(), 137 metadata: req.metadata, 138 creditor, 139 exchange_base_url: req.exchange_base_url, 140 }, 141 &Timestamp::now(), 142 ) 143 .await?; 144 match result { 145 db::TransferResult::Success { id, initiated_at } => Ok(TransferResponse { 146 timestamp: initiated_at.into(), 147 row_id: id, 148 }), 149 db::TransferResult::RequestUidReuse => { 150 Err(failure_code(ErrorCode::BANK_TRANSFER_REQUEST_UID_REUSED)) 151 } 152 db::TransferResult::WtidReuse => { 153 Err(failure_code(ErrorCode::BANK_TRANSFER_WTID_REUSED)) 154 } 155 } 156 } 157 158 async fn transfer_page( 159 &self, 160 page: Page, 161 status: Option<TransferState>, 162 ) -> ApiResult<TransferList> { 163 Ok(TransferList { 164 transfers: db::transfer_page(&self.pool, &status, &page).await?, 165 debit_account: self.payto.as_uri(), 166 }) 167 } 168 169 async fn transfer_by_id(&self, id: u64) -> ApiResult<Option<TransferStatus>> { 170 Ok(db::transfer_by_id(&self.pool, id).await?) 171 } 172 173 async fn outgoing_history(&self, params: History) -> ApiResult<OutgoingHistory> { 174 Ok(OutgoingHistory { 175 outgoing_transactions: db::outgoing_history(&self.pool, ¶ms, || { 176 self.taler_out_channel.subscribe() 177 }) 178 .await?, 179 debit_account: self.payto.as_uri(), 180 }) 181 } 182 183 async fn incoming_history(&self, params: History) -> ApiResult<IncomingHistory> { 184 Ok(IncomingHistory { 185 incoming_transactions: db::incoming_history(&self.pool, ¶ms, || { 186 self.taler_in_channel.subscribe() 187 }) 188 .await?, 189 credit_account: self.payto.as_uri(), 190 }) 191 } 192 193 async fn add_incoming_reserve( 194 &self, 195 req: AddIncomingRequest, 196 ) -> ApiResult<AddIncomingResponse> { 197 let debtor = FullHuPayto::try_from(&req.debit_account)?; 198 let res = db::register_tx_in_admin( 199 &self.pool, 200 &TxInAdmin { 201 amount: req.amount, 202 subject: format!("Admin incoming {}", req.reserve_pub), 203 debtor, 204 metadata: IncomingKey::reserve(req.reserve_pub), 205 }, 206 &Timestamp::now(), 207 ) 208 .await?; 209 match res { 210 AddIncomingResult::Success { 211 row_id, valued_at, .. 212 } => Ok(AddIncomingResponse { 213 row_id, 214 timestamp: date_to_utc_ts(&valued_at).into(), 215 }), 216 AddIncomingResult::ReservePubReuse => { 217 Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) 218 } 219 AddIncomingResult::UnknownMapping | AddIncomingResult::MappingReuse => { 220 unreachable!("mapping not used") 221 } 222 } 223 } 224 225 async fn add_incoming_kyc(&self, req: AddKycauthRequest) -> ApiResult<AddIncomingResponse> { 226 let debtor = FullHuPayto::try_from(&req.debit_account)?; 227 let res = db::register_tx_in_admin( 228 &self.pool, 229 &TxInAdmin { 230 amount: req.amount, 231 subject: format!("Admin incoming KYC:{}", req.account_pub), 232 debtor, 233 metadata: IncomingKey::kyc(req.account_pub), 234 }, 235 &Timestamp::now(), 236 ) 237 .await?; 238 match res { 239 AddIncomingResult::Success { 240 row_id, valued_at, .. 241 } => Ok(AddIncomingResponse { 242 row_id, 243 timestamp: date_to_utc_ts(&valued_at).into(), 244 }), 245 AddIncomingResult::ReservePubReuse => unreachable!("kyc"), 246 AddIncomingResult::UnknownMapping | AddIncomingResult::MappingReuse => { 247 unreachable!("mapping not used") 248 } 249 } 250 } 251 252 async fn add_incoming_mapped(&self, req: AddMappedRequest) -> ApiResult<AddIncomingResponse> { 253 let debtor = FullHuPayto::try_from(&req.debit_account)?; 254 let res = db::register_tx_in_admin( 255 &self.pool, 256 &TxInAdmin { 257 amount: req.amount, 258 subject: format!("Admin incoming MAP:{}", req.authorization_pub), 259 debtor, 260 metadata: IncomingKey::map(req.authorization_pub), 261 }, 262 &Timestamp::now(), 263 ) 264 .await?; 265 match res { 266 AddIncomingResult::Success { 267 row_id, valued_at, .. 268 } => Ok(AddIncomingResponse { 269 row_id, 270 timestamp: date_to_utc_ts(&valued_at).into(), 271 }), 272 AddIncomingResult::ReservePubReuse => { 273 Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) 274 } 275 AddIncomingResult::UnknownMapping => { 276 Err(failure_code(ErrorCode::BANK_TRANSFER_MAPPING_UNKNOWN)) 277 } 278 AddIncomingResult::MappingReuse => { 279 Err(failure_code(ErrorCode::BANK_TRANSFER_MAPPING_REUSED)) 280 } 281 } 282 } 283 284 fn support_account_check(&self) -> bool { 285 false 286 } 287 } 288 289 impl Revenue for MagnetApi { 290 async fn history(&self, params: History) -> ApiResult<RevenueIncomingHistory> { 291 Ok(RevenueIncomingHistory { 292 incoming_transactions: db::revenue_history(&self.pool, ¶ms, || { 293 self.in_channel.subscribe() 294 }) 295 .await?, 296 credit_account: self.payto.as_uri(), 297 }) 298 } 299 } 300 301 impl PreparedTransfer for MagnetApi { 302 fn supported_formats(&self) -> &[SubjectFormat] { 303 &[SubjectFormat::SIMPLE] 304 } 305 306 async fn registration(&self, req: RegistrationRequest) -> ApiResult<RegistrationResponse> { 307 let creditor = FullHuPayto::try_from(&req.credit_account)?; 308 if *creditor != *self.payto { 309 return Err(failure_code(ErrorCode::BANK_UNKNOWN_CREDITOR)); 310 } 311 match db::transfer_register(&self.pool, &req).await? { 312 db::RegistrationResult::Success => { 313 let simple = TransferSubject::Simple { 314 credit_amount: req.credit_amount, 315 subject: if req.authorization_pub == req.account_pub && !req.recurrent { 316 fmt_in_subject(req.r#type.into(), &req.account_pub).to_string() 317 } else { 318 fmt_in_subject(IncomingType::map, &req.authorization_pub).to_string() 319 }, 320 }; 321 ApiResult::Ok(RegistrationResponse { 322 subjects: vec![simple], 323 expiration: TalerTimestamp::Never, 324 }) 325 } 326 db::RegistrationResult::ReservePubReuse => { 327 ApiResult::Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) 328 } 329 } 330 } 331 332 async fn unregistration(&self, req: Unregistration) -> ApiResult<bool> { 333 Ok(db::transfer_unregister(&self.pool, &req).await?) 334 } 335 } 336 337 impl Observability for MagnetApi { 338 async fn metrics(&self) -> ApiResult<&Registry> { 339 self.metrics.sync(&self.pool).await; 340 Ok(&self.registry) 341 } 342 } 343 344 #[cfg(test)] 345 mod test { 346 347 use std::sync::{ 348 Arc, LazyLock, 349 atomic::{AtomicU64, Ordering}, 350 }; 351 352 use jiff::{Timestamp, Zoned}; 353 use sqlx::{PgPool, Row as _, postgres::PgRow}; 354 use taler_api::{ 355 api::TalerRouter as _, 356 auth::AuthMethod, 357 db::TypeHelper as _, 358 subject::{IncomingKey, OutgoingSubject}, 359 }; 360 use taler_common::{ 361 api::{ 362 EddsaPublicKey, 363 observability::Config, 364 prepared::PreparedTransferConfig, 365 revenue::RevenueConfig, 366 wire::{TransferState, WireConfig}, 367 }, 368 db::IncomingType, 369 types::{ 370 amount::amount, 371 payto::{PaytoURI, payto}, 372 }, 373 }; 374 use taler_macros::db_test; 375 use taler_test_utils::{ 376 Router, 377 routine::{ 378 Status, admin_add_incoming_routine, in_history_routine, out_history_routine, 379 registration_routine, revenue_routine, transfer_routine, 380 }, 381 server::TestServer, 382 tasks, 383 }; 384 385 use crate::{ 386 FullHuPayto, 387 api::MagnetApi, 388 db::{self, TxIn, TxOutKind}, 389 magnet_api::types::TxStatus, 390 magnet_payto, 391 }; 392 393 static PAYTO: LazyLock<FullHuPayto> = LazyLock::new(|| { 394 magnet_payto("payto://iban/HU02162000031000164800000000?receiver-name=Smith") 395 }); 396 static EXCHANGE: LazyLock<PaytoURI> = LazyLock::new(|| PAYTO.as_uri()); 397 static UNKNOWN: LazyLock<PaytoURI> = 398 LazyLock::new(|| payto("payto://iban/HU60162006491000639900000000?receiver-name=Unknown")); 399 400 async fn setup(pool: &PgPool) -> Router { 401 let api = Arc::new(MagnetApi::start(pool.clone(), PAYTO.clone()).await); 402 Router::new() 403 .wire_gateway(api.clone(), AuthMethod::None) 404 .prepared_transfer(api.clone()) 405 .revenue(api.clone(), AuthMethod::None) 406 .observability(api, AuthMethod::None) 407 .finalize() 408 } 409 410 #[db_test] 411 async fn config(pool: PgPool) { 412 let server = setup(&pool).await; 413 server 414 .get("/taler-wire-gateway/config") 415 .await 416 .assert_ok_json::<WireConfig>(); 417 server 418 .get("/taler-prepared-transfer/config") 419 .await 420 .assert_ok_json::<PreparedTransferConfig>(); 421 server 422 .get("/taler-revenue/config") 423 .await 424 .assert_ok_json::<RevenueConfig>(); 425 server 426 .get("/taler-observability/config") 427 .await 428 .assert_ok_json::<Config>(); 429 server.get("/taler-observability/metrics").await.assert_ok(); 430 } 431 432 #[db_test] 433 async fn transfer(pool: PgPool) { 434 let server = setup(&pool).await; 435 transfer_routine( 436 &server.prefix("/taler-wire-gateway"), 437 TransferState::pending, 438 &payto("payto://iban/HU02162000031000164800000000?receiver-name=name"), 439 ) 440 .await; 441 } 442 443 static CODE: AtomicU64 = AtomicU64::new(0); 444 445 async fn r#in(db: &PgPool, subject: Option<IncomingKey>) { 446 db::register_tx_in( 447 &mut db.acquire().await.unwrap(), 448 &TxIn { 449 code: CODE.fetch_add(1, Ordering::Relaxed), 450 amount: amount("EUR:10"), 451 subject: "subject".into(), 452 debtor: magnet_payto( 453 "payto://iban/HU30162000031000163100000000?receiver-name=name", 454 ), 455 value_date: Zoned::now().date(), 456 status: TxStatus::Completed, 457 }, 458 &subject, 459 &Timestamp::now(), 460 ) 461 .await 462 .unwrap(); 463 } 464 465 async fn in_malformed(db: &PgPool) { 466 r#in(db, None).await 467 } 468 469 async fn in_talerable(db: &PgPool) { 470 r#in(db, Some(IncomingKey::reserve(EddsaPublicKey::rand()))).await 471 } 472 473 async fn out(db: &PgPool, kind: &TxOutKind) { 474 db::register_tx_out( 475 &mut db.acquire().await.unwrap(), 476 &db::TxOut { 477 code: CODE.fetch_add(1, Ordering::Relaxed), 478 amount: amount("EUR:10"), 479 subject: "subject".into(), 480 creditor: PAYTO.clone(), 481 value_date: Zoned::now().date(), 482 status: TxStatus::Completed, 483 }, 484 kind, 485 &Timestamp::now(), 486 ) 487 .await 488 .unwrap(); 489 } 490 491 async fn out_talerable(db: &PgPool) { 492 out(db, &TxOutKind::Talerable(OutgoingSubject::rand())).await 493 } 494 495 async fn out_bounce(db: &PgPool) { 496 out(db, &TxOutKind::Bounce(CODE.load(Ordering::Relaxed) as u32)).await 497 } 498 499 async fn out_malformed(db: &PgPool) { 500 out(db, &TxOutKind::Simple).await 501 } 502 503 #[db_test] 504 async fn outgoing_history(db: PgPool) { 505 let db = &db; 506 let server = setup(db).await; 507 out_history_routine( 508 &server.prefix("/taler-wire-gateway"), 509 tasks!({ out_talerable(db).await }), 510 tasks!( 511 { out_bounce(db).await }, 512 { out_malformed(db).await }, 513 { in_malformed(db).await }, 514 { in_talerable(db).await } 515 ), 516 ) 517 .await; 518 } 519 520 #[db_test] 521 async fn admin_add_incoming(db: PgPool) { 522 let server = setup(&db).await; 523 admin_add_incoming_routine( 524 &server.prefix("/taler-wire-gateway"), 525 &server.prefix("/taler-prepared-transfer"), 526 &EXCHANGE, 527 &EXCHANGE, 528 ) 529 .await; 530 } 531 532 #[db_test] 533 async fn in_history(db: PgPool) { 534 let db = &db; 535 let server = setup(db).await; 536 in_history_routine( 537 &server.prefix("/taler-wire-gateway"), 538 &server.prefix("/taler-prepared-transfer"), 539 &EXCHANGE, 540 &EXCHANGE, 541 tasks!({ in_talerable(db).await }), 542 tasks!( 543 { out_malformed(db).await }, 544 { out_talerable(db).await }, 545 { out_bounce(db).await }, 546 { in_malformed(db).await } 547 ), 548 ) 549 .await; 550 } 551 552 #[db_test] 553 async fn revenue(db: PgPool) { 554 let db = &db; 555 let server = setup(db).await; 556 revenue_routine( 557 &server.prefix("/taler-wire-gateway"), 558 &server.prefix("/taler-revenue"), 559 &EXCHANGE, 560 tasks!({ in_malformed(db).await }, { in_talerable(db).await },), 561 tasks!({ out_malformed(db).await }, { out_talerable(db).await }, { 562 out_bounce(db).await 563 }), 564 ) 565 .await; 566 } 567 568 async fn check_in(pool: &PgPool) -> Vec<Status> { 569 sqlx::query( 570 " 571 SELECT pending_recurrent_in.authorization_pub IS NOT NULL, initiated_id IS NOT NULL, type, metadata 572 FROM tx_in 573 LEFT JOIN taler_in USING (tx_in_id) 574 LEFT JOIN pending_recurrent_in USING (tx_in_id) 575 LEFT JOIN bounced USING (tx_in_id) 576 ORDER BY tx_in.tx_in_id 577 ", 578 ) 579 .try_map(|r: PgRow| { 580 Ok( 581 if r.try_get_flag(0)? { 582 Status::Pending 583 } else if r.try_get_flag(1)? { 584 Status::Bounced 585 } else { 586 match r.try_get(2)? { 587 None => Status::Simple, 588 Some(IncomingType::reserve) => Status::Reserve(r.try_get(3)?), 589 Some(IncomingType::kyc) => Status::Kyc(r.try_get(3)?), 590 Some(e) => unreachable!("{e:?}") 591 } 592 } 593 ) 594 }) 595 .fetch_all(pool) 596 .await 597 .unwrap() 598 } 599 600 #[db_test] 601 async fn registration(db: PgPool) { 602 let server = setup(&db).await; 603 registration_routine( 604 &server.prefix("/taler-wire-gateway"), 605 &server.prefix("/taler-prepared-transfer"), 606 &EXCHANGE, 607 &EXCHANGE, 608 &UNKNOWN, 609 || check_in(&db), 610 ) 611 .await; 612 } 613 }