api.rs (17451B)
1 /* 2 * This file is part of LibEuFin. 3 * Copyright (C) 2026 Taler Systems S.A. 4 5 * LibEuFin is free software; you can redistribute it and/or modify 6 * it under the terms of the GNU Affero General Public License as 7 * published by the Free Software Foundation; either version 3, or 8 * (at your option) any later version. 9 10 * LibEuFin is distributed in the hope that it will be useful, but 11 * WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY 12 * or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Affero General 13 * Public License for more details. 14 15 * You should have received a copy of the GNU Affero General Public 16 * License along with LibEuFin; see the file COPYING. If not, see 17 * <http://www.gnu.org/licenses/> 18 */ 19 20 use const_format::formatcp; 21 use jiff::Timestamp; 22 use libeufin_ebics::{ 23 ebics::{TaskStatus, rand_ebics_id}, 24 iso20022::model::{InId, InTx}, 25 }; 26 use prometheus_client::{ 27 encoding::EncodeLabelSet, 28 metrics::{family::Family, gauge::Gauge}, 29 registry::{Registry, Unit}, 30 }; 31 use sqlx::{PgPool, types::Json}; 32 use taler_api::{ 33 api::{ 34 TalerApi, 35 observability::Observability, 36 prepared::{PreparedTransfer, simple_subject}, 37 revenue::Revenue, 38 wire::WireGateway, 39 }, 40 error::{ApiResult, failure_code}, 41 subject::{IncomingKey, fmt_in_subject, subject_fmt_qr_bill}, 42 }; 43 use taler_common::{ 44 api::{ 45 params::{History, Page}, 46 prepared::{ 47 RegistrationRequest, RegistrationResponse, SubjectFormat, TransferSubject, 48 Unregistration, 49 }, 50 revenue::RevenueIncomingHistory, 51 wire::{ 52 AddIncomingRequest, AddIncomingResponse, AddKycauthRequest, AddMappedRequest, 53 IncomingHistory, OutgoingHistory, TransferList, TransferRequest, TransferResponse, 54 TransferState, TransferStatus, 55 }, 56 }, 57 error_code::ErrorCode, 58 types::{ 59 amount::{Amount, Currency}, 60 iban::IBAN, 61 payto::{FullIbanPayto, IbanPayto, PaytoURI}, 62 time::TalerTimestamp, 63 }, 64 }; 65 use tokio::sync::watch::Sender; 66 67 use crate::{ 68 constants::{FETCH_TASK_KEY, SUBMIT_TASK_KEY}, 69 db::{ 70 self, 71 exchange::{ 72 TransferResult, incoming_history, outgoing_history, revenue_history, transfer, 73 transfer_by_id, transfer_page, 74 }, 75 payment::{IncomingRegistrationResult, register_in_talerable}, 76 transfer::{RegistrationResult, transfer_register, transfer_unregister}, 77 }, 78 }; 79 80 pub struct NexusApi { 81 pub pool: sqlx::PgPool, 82 pub currency: Currency, 83 pub payto: FullIbanPayto, 84 pub qr_iban: Option<IBAN>, 85 pub in_channel: Sender<i64>, 86 pub taler_in_channel: Sender<i64>, 87 pub taler_out_channel: Sender<i64>, 88 metrics: Metrics, 89 registry: Registry, 90 } 91 92 #[derive(Clone, Debug, Hash, PartialEq, Eq, EncodeLabelSet)] 93 struct TaskLabel { 94 name: &'static str, 95 } 96 97 #[derive(Default)] 98 pub struct Metrics { 99 db_access: Gauge, 100 task_execution: Family<TaskLabel, Gauge>, 101 task_success: Family<TaskLabel, Gauge>, 102 } 103 104 impl Metrics { 105 pub fn registry(&self) -> Registry { 106 let mut registry = Registry::default(); 107 108 registry.register( 109 "db_access", 110 "Whether the last database metrics refresh succeeded", 111 self.db_access.clone(), 112 ); 113 114 registry.register_with_unit( 115 "task_execution_timestamp_seconds", 116 "Unix timestamp of the last task execution", 117 Unit::Seconds, 118 self.task_execution.clone(), 119 ); 120 121 registry.register_with_unit( 122 "task_success_timestamp_seconds", 123 "Unix timestamp of the last successful task execution", 124 Unit::Seconds, 125 self.task_success.clone(), 126 ); 127 128 registry 129 } 130 131 pub async fn sync(&self, db: &PgPool) { 132 let success = 133 match sqlx::query_as::<_, (Option<Json<TaskStatus>>, Option<Json<TaskStatus>>)>( 134 formatcp!( 135 " 136 SELECT 137 (SELECT value FROM kv WHERE key='{SUBMIT_TASK_KEY}'), 138 (SELECT value FROM kv WHERE key='{FETCH_TASK_KEY}') 139 " 140 ), 141 ) 142 .fetch_one(db) 143 .await 144 { 145 Ok((submit, fetch)) => { 146 for (name, status) in [("submit", submit), ("fetch", fetch)] { 147 if let Some(status) = status { 148 if let Some(time) = status.last_trial { 149 self.task_execution 150 .get_or_create(&TaskLabel { name }) 151 .set(time.as_second()); 152 } 153 if let Some(time) = status.last_successfull { 154 self.task_success 155 .get_or_create(&TaskLabel { name }) 156 .set(time.as_second()); 157 } 158 } 159 } 160 true 161 } 162 Err(_) => false, 163 }; 164 self.db_access.set(if success { 1 } else { 0 }); 165 } 166 } 167 168 impl NexusApi { 169 pub async fn start( 170 pool: sqlx::PgPool, 171 currency: Currency, 172 payto: FullIbanPayto, 173 qr_iban: Option<IBAN>, 174 ) -> Self { 175 let in_channel = Sender::new(0); 176 let taler_in_channel = Sender::new(0); 177 let taler_out_channel = Sender::new(0); 178 179 let metrics = Metrics::default(); 180 181 let tmp = Self { 182 pool: pool.clone(), 183 payto, 184 qr_iban, 185 currency, 186 in_channel: in_channel.clone(), 187 taler_in_channel: taler_in_channel.clone(), 188 taler_out_channel: taler_out_channel.clone(), 189 registry: metrics.registry(), 190 metrics, 191 }; 192 tokio::spawn(db::notification_listener( 193 pool, 194 in_channel, 195 taler_in_channel, 196 taler_out_channel, 197 )); 198 tmp 199 } 200 } 201 202 impl TalerApi for NexusApi { 203 fn currency(&self) -> Currency { 204 self.currency 205 } 206 207 fn implementation(&self) -> &'static str { 208 "urn:net:taler:specs:libeufin-nexus:taler-rust" 209 } 210 } 211 212 async fn add_incoming( 213 db: &PgPool, 214 subject: &IncomingKey, 215 amount: Amount, 216 debit_account: PaytoURI, 217 ) -> ApiResult<AddIncomingResponse> { 218 FullIbanPayto::try_from(&debit_account)?; 219 let now = Timestamp::now(); 220 match register_in_talerable( 221 db, 222 &InTx { 223 id: InId { 224 uetr: None, 225 tx_id: Some(rand_ebics_id()), 226 sref: None, 227 }, 228 amount, 229 credit_fee: Amount::zero(&amount.currency), 230 subject: Some(format!( 231 "Manual incoming {}", 232 fmt_in_subject(subject.ty, &subject.key) 233 )), 234 execution_time: now, 235 debtor: Some(debit_account), 236 }, 237 subject, 238 ) 239 .await? 240 { 241 IncomingRegistrationResult::Success(in_result) => Ok(AddIncomingResponse { 242 row_id: in_result.id, 243 timestamp: now.into(), 244 }), 245 IncomingRegistrationResult::ReservePubReuse => { 246 Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) 247 } 248 IncomingRegistrationResult::MappingReuse => { 249 Err(failure_code(ErrorCode::BANK_TRANSFER_MAPPING_REUSED)) 250 } 251 IncomingRegistrationResult::UnknownMapping => { 252 Err(failure_code(ErrorCode::BANK_TRANSFER_MAPPING_UNKNOWN)) 253 } 254 } 255 } 256 257 impl WireGateway for NexusApi { 258 async fn transfer(&self, req: TransferRequest) -> ApiResult<TransferResponse> { 259 FullIbanPayto::try_from(&req.credit_account)?; 260 let result = transfer(&self.pool, &req, &rand_ebics_id(), &Timestamp::now()).await?; 261 match result { 262 TransferResult::Success { id, timestamp } => Ok(TransferResponse { 263 timestamp: timestamp.into(), 264 row_id: id, 265 }), 266 TransferResult::RequestUidReuse => { 267 Err(failure_code(ErrorCode::BANK_TRANSFER_REQUEST_UID_REUSED)) 268 } 269 TransferResult::WtidReuse => Err(failure_code(ErrorCode::BANK_TRANSFER_WTID_REUSED)), 270 } 271 } 272 273 async fn transfer_page( 274 &self, 275 page: Page, 276 status: Option<TransferState>, 277 ) -> ApiResult<TransferList> { 278 Ok(TransferList { 279 transfers: transfer_page(&self.pool, &self.currency, &page, &status).await?, 280 debit_account: self.payto.as_uri(), 281 }) 282 } 283 284 async fn transfer_by_id(&self, id: u64) -> ApiResult<Option<TransferStatus>> { 285 Ok(transfer_by_id(&self.pool, &self.currency, id).await?) 286 } 287 288 async fn outgoing_history(&self, params: History) -> ApiResult<OutgoingHistory> { 289 Ok(OutgoingHistory { 290 outgoing_transactions: outgoing_history(&self.pool, &self.currency, ¶ms, || { 291 self.taler_out_channel.subscribe() 292 }) 293 .await?, 294 debit_account: self.payto.as_uri(), 295 }) 296 } 297 298 async fn incoming_history(&self, params: History) -> ApiResult<IncomingHistory> { 299 Ok(IncomingHistory { 300 incoming_transactions: incoming_history(&self.pool, &self.currency, ¶ms, || { 301 self.taler_in_channel.subscribe() 302 }) 303 .await?, 304 credit_account: self.payto.as_uri(), 305 }) 306 } 307 308 async fn add_incoming_reserve( 309 &self, 310 req: AddIncomingRequest, 311 ) -> ApiResult<AddIncomingResponse> { 312 add_incoming( 313 &self.pool, 314 &IncomingKey::reserve(req.reserve_pub), 315 req.amount, 316 req.debit_account, 317 ) 318 .await 319 } 320 321 async fn add_incoming_kyc(&self, req: AddKycauthRequest) -> ApiResult<AddIncomingResponse> { 322 add_incoming( 323 &self.pool, 324 &IncomingKey::kyc(req.account_pub), 325 req.amount, 326 req.debit_account, 327 ) 328 .await 329 } 330 331 async fn add_incoming_mapped(&self, req: AddMappedRequest) -> ApiResult<AddIncomingResponse> { 332 add_incoming( 333 &self.pool, 334 &IncomingKey::map(req.authorization_pub), 335 req.amount, 336 req.debit_account, 337 ) 338 .await 339 } 340 341 fn support_account_check(&self) -> bool { 342 false 343 } 344 } 345 346 impl Revenue for NexusApi { 347 async fn history(&self, params: History) -> ApiResult<RevenueIncomingHistory> { 348 Ok(RevenueIncomingHistory { 349 incoming_transactions: revenue_history(&self.pool, &self.currency, ¶ms, || { 350 self.in_channel.subscribe() 351 }) 352 .await?, 353 credit_account: self.payto.as_uri(), 354 }) 355 } 356 } 357 358 impl PreparedTransfer for NexusApi { 359 fn supported_formats(&self) -> &[SubjectFormat] { 360 if self.qr_iban.is_some() { 361 &[SubjectFormat::SIMPLE, SubjectFormat::CH_QR_BILL] 362 } else { 363 &[SubjectFormat::SIMPLE] 364 } 365 } 366 367 async fn registration(&self, req: RegistrationRequest) -> ApiResult<RegistrationResponse> { 368 let creditor = IbanPayto::try_from(&req.credit_account)?; 369 let reference_number = if creditor.iban == self.payto.iban { 370 None 371 } else if Some(creditor.iban) == self.qr_iban { 372 Some(subject_fmt_qr_bill(req.authorization_pub.as_ref())) 373 } else { 374 return Err(failure_code(ErrorCode::BANK_UNKNOWN_CREDITOR)); 375 }; 376 match transfer_register( 377 &self.pool, 378 req.r#type.into(), 379 &req.account_pub, 380 &req.authorization_pub, 381 &req.authorization_sig, 382 req.recurrent, 383 reference_number.as_deref(), 384 &Timestamp::now(), 385 ) 386 .await? 387 { 388 RegistrationResult::Success => ApiResult::Ok(RegistrationResponse { 389 subjects: vec![if let Some(qr_reference_number) = reference_number { 390 TransferSubject::QrBill { 391 credit_amount: req.credit_amount, 392 qr_reference_number, 393 } 394 } else { 395 simple_subject(req) 396 }], 397 expiration: TalerTimestamp::Never, 398 }), 399 RegistrationResult::ReservePubReuse => { 400 ApiResult::Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT)) 401 } 402 RegistrationResult::SubjectReuse => { 403 ApiResult::Err(failure_code(ErrorCode::BANK_DERIVATION_REUSE)) 404 } 405 } 406 } 407 408 async fn unregistration(&self, req: Unregistration) -> ApiResult<bool> { 409 Ok(transfer_unregister(&self.pool, &req.authorization_pub, &Timestamp::now()).await?) 410 } 411 } 412 413 impl Observability for NexusApi { 414 async fn metrics(&self) -> ApiResult<&Registry> { 415 self.metrics.sync(&self.pool).await; 416 Ok(&self.registry) 417 } 418 } 419 420 #[cfg(test)] 421 pub mod test { 422 use std::sync::Arc; 423 424 use sqlx::PgPool; 425 use taler_api::{api::TalerRouter as _, auth::AuthMethod, subject::OutgoingSubject}; 426 use taler_common::api::{ 427 prepared::PreparedTransferConfig, 428 revenue::RevenueConfig, 429 wire::{TransferState, WireConfig}, 430 }; 431 use taler_macros::db_test; 432 use taler_test_utils::{ 433 Router, 434 routine::{ 435 admin_add_incoming_routine, in_history_routine, out_history_routine, 436 registration_routine, revenue_routine, transfer_routine, 437 }, 438 server::TestServer as _, 439 tasks, 440 }; 441 442 use crate::{ 443 api::NexusApi, 444 db::{payment::register_out_tx, test::check_in}, 445 test::{CLIENT, CURR, EXCHANGE, UNKNOWN, gen_out_pay}, 446 }; 447 448 pub async fn api_setup(db: &PgPool) -> Router { 449 let api = Arc::new(NexusApi::start(db.clone(), CURR, EXCHANGE.clone(), None).await); 450 Router::new() 451 .wire_gateway(api.clone(), AuthMethod::None) 452 .prepared_transfer(api.clone()) 453 .revenue(api, AuthMethod::None) 454 .finalize() 455 } 456 457 #[db_test] 458 async fn config(db: PgPool) { 459 let server = api_setup(&db).await; 460 server 461 .get("/taler-wire-gateway/config") 462 .await 463 .assert_ok_json::<WireConfig>(); 464 server 465 .get("/taler-prepared-transfer/config") 466 .await 467 .assert_ok_json::<PreparedTransferConfig>(); 468 server 469 .get("/taler-revenue/config") 470 .await 471 .assert_ok_json::<RevenueConfig>(); 472 } 473 474 #[db_test] 475 async fn transfer(db: PgPool) { 476 let server = api_setup(&db).await; 477 transfer_routine( 478 &server.prefix("/taler-wire-gateway"), 479 TransferState::pending, 480 &EXCHANGE.as_uri(), 481 ) 482 .await; 483 // TODO 484 /*db.initiated.batchSubmissionSuccess(1, Instant.now(), "ORDER1") 485 db.initiated.batchSubmissionFailure(2, Instant.now(), "Failure") 486 db.initiated.batchSubmissionFailure(3, Instant.now(), "Failure") 487 client.getA("/taler-wire-gateway/transfers?status=transient_failure").assertOkJson<TransferList> { 488 assertEquals(2, it.transfers.size) 489 } 490 client.getA("/taler-wire-gateway/transfers?status=pending").assertOkJson<TransferList> { 491 assertEquals(4, it.transfers.size) 492 }*/ 493 } 494 495 #[db_test] 496 async fn outgoing_history(db: PgPool) { 497 let server = api_setup(&db).await; 498 out_history_routine( 499 &server.prefix("/taler-wire-gateway"), 500 tasks!({ 501 register_out_tx(&db, &gen_out_pay("subject"), Some(&OutgoingSubject::rand())) 502 .await 503 .unwrap(); 504 }), 505 tasks!(), 506 ) 507 .await; 508 } 509 510 #[db_test] 511 async fn incoming_history(db: PgPool) { 512 let server = api_setup(&db).await; 513 in_history_routine( 514 &server.prefix("/taler-wire-gateway"), 515 &server.prefix("/taler-prepared-transfer"), 516 &CLIENT.as_uri(), 517 &EXCHANGE.as_uri(), 518 tasks!(), 519 tasks!(), 520 ) 521 .await; 522 } 523 524 #[db_test] 525 async fn admin_add_incoming(db: PgPool) { 526 let server = api_setup(&db).await; 527 admin_add_incoming_routine( 528 &server.prefix("/taler-wire-gateway"), 529 &server.prefix("/taler-prepared-transfer"), 530 &CLIENT.as_uri(), 531 &EXCHANGE.as_uri(), 532 ) 533 .await; 534 } 535 536 #[db_test] 537 async fn revenue(db: PgPool) { 538 let server = api_setup(&db).await; 539 revenue_routine( 540 &server.prefix("/taler-wire-gateway"), 541 &server.prefix("/taler-revenue"), 542 &CLIENT.as_uri(), 543 tasks!(), 544 tasks!(), 545 ) 546 .await; 547 } 548 549 #[db_test] 550 async fn registration(db: PgPool) { 551 let server = api_setup(&db).await; 552 registration_routine( 553 &server.prefix("/taler-wire-gateway"), 554 &server.prefix("/taler-prepared-transfer"), 555 &CLIENT.as_uri(), 556 &EXCHANGE.as_uri(), 557 &UNKNOWN, 558 || check_in(&db), 559 ) 560 .await; 561 } 562 }