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 }