taler-rust

GNU Taler code in Rust. Largely core banking integrations.
Log | Files | Refs | Submodules | README | LICENSE

sse.rs (2429B)


      1 /*
      2   This file is part of TALER
      3   Copyright (C) 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::{borrow::Cow, num::NonZeroUsize};
     18 
     19 use futures_util::StreamExt as _;
     20 use http_body_util::BodyDataStream;
     21 use hyper::body::Incoming;
     22 use sse_core::{SseDecoder, SseEvent, SseStream, SseStreamError};
     23 
     24 #[derive(Debug, Default, PartialEq, Eq)]
     25 pub struct SseMessage {
     26     pub event: Cow<'static, str>,
     27     pub data: String,
     28 }
     29 
     30 /// Server-sent event client
     31 pub struct SseClient {
     32     pub reconnection_time: Option<u64>,
     33     stream: sse_core::SseStream<BodyDataStream<Incoming>>,
     34 }
     35 
     36 impl SseClient {
     37     pub fn new() -> Self {
     38         Self {
     39             reconnection_time: None,
     40             stream: SseStream::with_decoder(SseDecoder::with_limit(
     41                 NonZeroUsize::new(1024 * 1024).unwrap(),
     42             )),
     43         }
     44     }
     45 
     46     pub fn last_event_id(&self) -> Option<&str> {
     47         self.stream.last_event_id().map(|it| it.as_ref())
     48     }
     49 
     50     pub fn connect(&mut self, body: Incoming) {
     51         self.stream.attach(BodyDataStream::new(body));
     52     }
     53 
     54     pub async fn next(&mut self) -> Option<Result<SseMessage, SseStreamError<hyper::Error>>> {
     55         loop {
     56             match self.stream.next().await {
     57                 Some(Ok(SseEvent::Message(message))) => {
     58                     return Some(Ok(SseMessage {
     59                         event: message.event,
     60                         data: message.data,
     61                     }));
     62                 }
     63                 Some(Ok(SseEvent::Retry(milliseconds))) => {
     64                     self.reconnection_time = Some(u64::from(milliseconds));
     65                 }
     66                 Some(Err(e)) => return Some(Err(e)),
     67                 None => return None,
     68             }
     69         }
     70     }
     71 }
     72 
     73 impl Default for SseClient {
     74     fn default() -> Self {
     75         Self::new()
     76     }
     77 }