notification.rs (2983B)
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 http_client::sse::SseClient; 18 use jiff::Timestamp; 19 use taler_common::ExpoBackoffDecorr; 20 use tokio::sync::Notify; 21 use tracing::{debug, error, trace}; 22 23 use crate::cyclos_api::{ 24 client::Client, 25 types::{NotificationEntityType, NotificationStatus}, 26 }; 27 28 pub async fn watch_notification(client: &Client<'_>, notify: &Notify) -> () { 29 let client_id = Timestamp::now().as_microsecond(); 30 let mut sse_client = SseClient::new(); 31 let mut jitter = ExpoBackoffDecorr::default(); 32 loop { 33 let res: anyhow::Result<()> = async { 34 loop { 35 // Register listener 36 let connected = client 37 .push_notifications(client_id, &mut sse_client) 38 .await?; 39 if !connected { 40 return Ok(()) 41 } 42 jitter.reset(); 43 // Read available ones 44 while let Some(message) = sse_client.next().await { 45 let msg = message?; 46 trace!(target: "notification", "new message {}: {}", msg.event, msg.data); 47 if msg.event == "newNotification" { 48 let deserializer = &mut serde_json::Deserializer::from_str(&msg.data); 49 let status: NotificationStatus = 50 serde_path_to_error::deserialize(deserializer)?; 51 debug!(target: "notification", "new notification {} {:?} {:?}", status.notification.id, status.notification.ty, status.notification.entity_type); 52 if status.notification.entity_type == Some(NotificationEntityType::Transfer) 53 { 54 notify.notify_waiters(); 55 } 56 // Find a way to buffer all transactions 57 } 58 } 59 tokio::time::sleep( jitter.backoff()).await; 60 } 61 } 62 .await; 63 if let Err(err) = res { 64 error!(target: "notification", "{err}"); 65 tokio::time::sleep(jitter.backoff()).await; 66 } else { 67 error!(target: "notification", "SSE connection refused with HTTP 204"); 68 std::future::pending::<()>().await; 69 } 70 } 71 }