leptos_use/use_event_source.rs
1use crate::ReconnectLimit;
2use crate::core::ConnectionReadyState;
3use codee::Decoder;
4use default_struct_builder::DefaultBuilder;
5use leptos::prelude::*;
6use std::fmt::Debug;
7use std::marker::PhantomData;
8use std::sync::Arc;
9use thiserror::Error;
10use wasm_bindgen::JsCast;
11
12/// Reactive [EventSource](https://developer.mozilla.org/en-US/docs/Web/API/EventSource)
13///
14/// An [EventSource](https://developer.mozilla.org/en-US/docs/Web/API/EventSource) or
15/// [Server-Sent-Events](https://developer.mozilla.org/en-US/docs/Web/API/Server-sent_events)
16/// instance opens a persistent connection to an HTTP server,
17/// which sends events in text/event-stream format.
18///
19/// ## Usage
20///
21/// Values are decoded via the given decoder. You can use any of the string codecs or a
22/// binary codec wrapped in `Base64`.
23///
24/// > Please check [the codec chapter](https://leptos-use.rs/codecs.html) to see what codecs are
25/// > available and what feature flags they require.
26///
27/// ```
28/// # use leptos::prelude::*;
29/// # use leptos_use::{use_event_source, UseEventSourceReturn};
30/// # use codee::string::JsonSerdeCodec;
31/// # use serde::{Deserialize, Serialize};
32/// #
33/// #[derive(Serialize, Deserialize, Clone, PartialEq)]
34/// pub struct EventSourceData {
35/// pub message: String,
36/// pub priority: u8,
37/// }
38///
39/// # #[component]
40/// # fn Demo() -> impl IntoView {
41/// let UseEventSourceReturn {
42/// ready_state, message, error, close, ..
43/// } = use_event_source::<EventSourceData, JsonSerdeCodec>("https://event-source-url");
44/// #
45/// # view! { }
46/// # }
47/// ```
48///
49/// ### Named Events
50///
51/// You can define named events when using `use_event_source_with_options`.
52///
53/// ```
54/// # use leptos::prelude::*;
55/// # use leptos_use::{use_event_source_with_options, UseEventSourceReturn, UseEventSourceOptions};
56/// # use codee::string::FromToStringCodec;
57/// #
58/// # #[component]
59/// # fn Demo() -> impl IntoView {
60/// let UseEventSourceReturn {
61/// ready_state, message, error, close, ..
62/// } = use_event_source_with_options::<String, FromToStringCodec>(
63/// "https://event-source-url",
64/// UseEventSourceOptions::default()
65/// .named_events(["notice".to_string(), "update".to_string()])
66/// );
67/// #
68/// # view! { }
69/// # }
70/// ```
71///
72/// ### Custom Event Handler
73///
74/// You can provide a custom `on_event` handler using `use_event_source_with_options`.
75/// `on_event` wil be run for every received event, including the built-in `open`, `error`,
76/// and `message` events, as well as any named events you have specified.
77///
78/// With the return value of `on_event` you can control, whether `message` and named events
79/// should be further processed by `use_event_source` (`UseEventSourceOnEventReturn::ProcessMessage`)
80/// or ignored (`UseEventSourceOnEventReturn::IgnoreProcessingMessage`).
81///
82/// By default, the handler returns `UseEventSourceOnEventReturn::ProcessMessage`.
83///
84/// ```
85/// # use leptos::prelude::*;
86/// # use leptos_use::{use_event_source_with_options, UseEventSourceReturn, UseEventSourceOptions, UseEventSourceMessage, UseEventSourceOnEventReturn};
87/// # use codee::string::FromToStringCodec;
88/// #
89/// # #[component]
90/// # fn Demo() -> impl IntoView {
91/// // Custom example handler: log event name and check for named `custom_error` event
92/// let custom_event_handler = |e: &web_sys::Event| {
93/// leptos::logging::log!("Received event: {}", e.type_());
94/// if e.type_() == "custom_error" {
95/// if let Ok(error_message) = UseEventSourceMessage::<String, FromToStringCodec>::try_from(e.clone()) {
96/// // Decoded successfully, log the error message
97/// leptos::logging::log!("Error message: {}", error_message.data);
98/// // skip processing this message event further
99/// return UseEventSourceOnEventReturn::IgnoreProcessingMessage;
100/// }
101/// }
102/// // Process other message events normally
103/// UseEventSourceOnEventReturn::ProcessMessage
104/// };
105/// let UseEventSourceReturn {
106/// ready_state, message, error, close, ..
107/// } = use_event_source_with_options::<String, FromToStringCodec>(
108/// "https://event-source-url",
109/// UseEventSourceOptions::default()
110/// .named_events(["custom_error".to_string()])
111/// .on_event(custom_event_handler)
112/// );
113/// #
114/// # view! { }
115/// # }
116/// ```
117///
118/// ### Immediate
119///
120/// Auto-connect (enabled by default).
121///
122/// This will call `open()` automatically for you, and you don't need to call it by yourself.
123///
124/// ### Auto-Reconnection
125///
126/// Reconnect on errors automatically (enabled by default).
127///
128/// You can control the number of reconnection attempts by setting `reconnect_limit` and the
129/// interval between them by setting `reconnect_interval`.
130///
131/// ```
132/// # use leptos::prelude::*;
133/// # use leptos_use::{use_event_source_with_options, UseEventSourceReturn, UseEventSourceOptions, ReconnectLimit};
134/// # use codee::string::FromToStringCodec;
135/// #
136/// # #[component]
137/// # fn Demo() -> impl IntoView {
138/// let UseEventSourceReturn {
139/// ready_state, message, error, close, ..
140/// } = use_event_source_with_options::<bool, FromToStringCodec>(
141/// "https://event-source-url",
142/// UseEventSourceOptions::default()
143/// .reconnect_limit(ReconnectLimit::Limited(5)) // at most 5 attempts
144/// .reconnect_interval(2000) // wait for 2 seconds between attempts
145/// );
146/// #
147/// # view! { }
148/// # }
149/// ```
150///
151///
152/// ## SendWrapped Return
153///
154/// The returned closures `open` and `close` are sendwrapped functions. They can
155/// only be called from the same thread that called `use_event_source`.
156///
157/// To disable auto-reconnection, set `reconnect_limit` to `0`.
158///
159/// ## Server-Side Rendering
160///
161/// > Make sure you follow the [instructions in Server-Side Rendering](https://leptos-use.rs/server_side_rendering.html).
162///
163/// On the server-side, `use_event_source` will always return `ready_state` as `ConnectionReadyState::Closed`,
164/// `data`, `event` and `error` will always be `None`, and `open` and `close` will do nothing.
165pub fn use_event_source<T, C>(
166 url: impl Into<Signal<String>>,
167) -> UseEventSourceReturn<
168 T,
169 C,
170 C::Error,
171 impl Fn() + Clone + Send + Sync + 'static,
172 impl Fn() + Clone + Send + Sync + 'static,
173>
174where
175 T: Clone + PartialEq + Send + Sync + 'static,
176 C: Decoder<T, Encoded = str> + Send + Sync,
177 C::Error: Send + Sync,
178{
179 use_event_source_with_options::<T, C>(url, UseEventSourceOptions::<T>::default())
180}
181
182/// Version of [`use_event_source`] that takes a `UseEventSourceOptions`. See [`use_event_source`] for how to use.
183pub fn use_event_source_with_options<T, C>(
184 url: impl Into<Signal<String>>,
185 options: UseEventSourceOptions<T>,
186) -> UseEventSourceReturn<
187 T,
188 C,
189 C::Error,
190 impl Fn() + Clone + Send + Sync + 'static,
191 impl Fn() + Clone + Send + Sync + 'static,
192>
193where
194 T: Clone + PartialEq + Send + Sync + 'static,
195 C: Decoder<T, Encoded = str> + Send + Sync,
196 C::Error: Send + Sync,
197{
198 let UseEventSourceOptions {
199 reconnect_limit,
200 reconnect_interval,
201 on_failed,
202 immediate,
203 named_events,
204 on_event,
205 with_credentials,
206 _marker,
207 } = options;
208
209 let (message, set_message) = signal(None::<UseEventSourceMessage<T, C>>);
210 let (ready_state, set_ready_state) = signal(ConnectionReadyState::Closed);
211 let (error, set_error) = signal(None::<UseEventSourceError<C::Error>>);
212
213 let open;
214 let close;
215
216 #[cfg(not(feature = "ssr"))]
217 {
218 use crate::{sendwrap_fn, use_event_listener};
219 use leptos::leptos_dom::helpers::TimeoutHandle;
220 use std::sync::atomic::{AtomicBool, AtomicU32};
221 use std::time::Duration;
222 use wasm_bindgen::prelude::*;
223
224 let (event_source, set_event_source) = signal_local(None::<web_sys::EventSource>);
225 let explicitly_closed = Arc::new(AtomicBool::new(false));
226 let retried = Arc::new(AtomicU32::new(0));
227
228 let on_event_return = move |e: &web_sys::Event| {
229 // make sure handler doesn't create reactive dependencies
230 #[cfg(debug_assertions)]
231 let _ = leptos::reactive::diagnostics::SpecialNonReactiveZone::enter();
232
233 on_event(e)
234 };
235
236 let on_message_event = {
237 let on_event_return = on_event_return.clone();
238 move |e: &web_sys::Event| {
239 match on_event_return(e) {
240 UseEventSourceOnEventReturn::IgnoreProcessingMessage => {
241 // skip processing message event!
242 }
243 UseEventSourceOnEventReturn::ProcessMessage => {
244 let message_event = e
245 .dyn_ref::<web_sys::MessageEvent>()
246 .expect("Event is not a MessageEvent");
247
248 match UseEventSourceMessage::<T, C>::try_from(message_event) {
249 Ok(event_msg) => {
250 set_message.set(Some(event_msg));
251 }
252 Err(err) => {
253 set_error.set(Some(err));
254 }
255 }
256 }
257 }
258 }
259 };
260
261 let init = StoredValue::new(None::<Arc<dyn Fn() + Send + Sync>>);
262
263 let reconnect_timer: StoredValue<Option<TimeoutHandle>> = StoredValue::new(None);
264
265 let clear_reconnect_timer = move || {
266 reconnect_timer.update_value(|timer| {
267 if let Some(timer) = timer.take() {
268 timer.clear();
269 }
270 });
271 };
272
273 let set_init = {
274 let explicitly_closed = Arc::clone(&explicitly_closed);
275 let retried = Arc::clone(&retried);
276
277 move |url: String| {
278 init.set_value(Some(Arc::new({
279 let explicitly_closed = Arc::clone(&explicitly_closed);
280 let retried = Arc::clone(&retried);
281 let on_event_return = on_event_return.clone();
282 let on_message_event = on_message_event.clone();
283 let named_events = named_events.clone();
284 let on_failed = Arc::clone(&on_failed);
285
286 move || {
287 if explicitly_closed.load(std::sync::atomic::Ordering::Relaxed) {
288 return;
289 }
290
291 let event_src_opts = web_sys::EventSourceInit::new();
292 event_src_opts.set_with_credentials(with_credentials);
293
294 let es = web_sys::EventSource::new_with_event_source_init_dict(
295 &url,
296 &event_src_opts,
297 )
298 .unwrap_throw();
299
300 set_ready_state.set(ConnectionReadyState::Connecting);
301
302 set_event_source.set(Some(es.clone()));
303
304 let on_open = Closure::wrap(Box::new({
305 let on_event_return = on_event_return.clone();
306 move |e: web_sys::Event| {
307 on_event_return(&e);
308 set_ready_state.set(ConnectionReadyState::Open);
309 set_error.set(None);
310 }})
311 as Box<dyn FnMut(web_sys::Event)>);
312 es.set_onopen(Some(on_open.as_ref().unchecked_ref()));
313 on_open.forget();
314
315 let on_error = Closure::wrap(Box::new({
316 let on_event_return = on_event_return.clone();
317 let explicitly_closed = Arc::clone(&explicitly_closed);
318 let retried = Arc::clone(&retried);
319 let on_failed = Arc::clone(&on_failed);
320 let es = es.clone();
321
322 move |e: web_sys::Event| {
323 on_event_return(&e);
324 set_ready_state.set(ConnectionReadyState::Closed);
325 set_error.set(Some(UseEventSourceError::ErrorEvent));
326
327 // only reconnect if EventSource isn't reconnecting by itself
328 // this is the case when the connection is closed (readyState is 2)
329 if es.ready_state() == 2
330 && !explicitly_closed.load(std::sync::atomic::Ordering::Relaxed)
331 {
332 es.close();
333
334 let retried_value = retried
335 .fetch_add(1, std::sync::atomic::Ordering::Relaxed)
336 + 1;
337
338 if !reconnect_limit.is_exceeded_by(retried_value as u64) {
339 clear_reconnect_timer();
340
341 reconnect_timer.set_value(
342 set_timeout_with_handle(
343 move || {
344 reconnect_timer.set_value(None);
345
346 if let Some(init) = init.get_value() {
347 init();
348 }
349 },
350 Duration::from_millis(reconnect_interval),
351 )
352 .ok(),
353 );
354 } else {
355 #[cfg(debug_assertions)]
356 let _z =
357 leptos::reactive::diagnostics::SpecialNonReactiveZone::enter();
358
359 on_failed();
360 }
361 }
362 }
363 })
364 as Box<dyn FnMut(web_sys::Event)>);
365 es.set_onerror(Some(on_error.as_ref().unchecked_ref()));
366 on_error.forget();
367
368 let on_message = Closure::wrap(Box::new({
369 let on_message_event = on_message_event.clone();
370 move |e: web_sys::MessageEvent| {
371 let e: &web_sys::Event = e.as_ref();
372 on_message_event(e);
373 }})
374 as Box<dyn FnMut(web_sys::MessageEvent)>);
375 es.set_onmessage(Some(on_message.as_ref().unchecked_ref()));
376 on_message.forget();
377
378 for event_name in named_events.clone() {
379 let event_handler = {
380 let on_message_event = on_message_event.clone();
381 move |e: web_sys::Event| {
382 on_message_event(&e);
383 }
384 };
385
386 let _ = use_event_listener(
387 es.clone(),
388 leptos::ev::Custom::<leptos::ev::Event>::new(event_name),
389 event_handler,
390 );
391 }
392 }
393 })))
394 }
395 };
396
397 close = {
398 let explicitly_closed = Arc::clone(&explicitly_closed);
399
400 sendwrap_fn!(move || {
401 // A reconnect scheduled by an earlier error would otherwise still fire and
402 // open a second connection behind the caller's back.
403 clear_reconnect_timer();
404
405 if let Some(event_source) = event_source.get_untracked() {
406 event_source.close();
407 set_event_source.set(None);
408 set_ready_state.set(ConnectionReadyState::Closed);
409 explicitly_closed.store(true, std::sync::atomic::Ordering::Relaxed);
410 }
411 })
412 };
413
414 let url: Signal<String> = url.into();
415
416 open = {
417 let close = close.clone();
418 let explicitly_closed = Arc::clone(&explicitly_closed);
419 let retried = Arc::clone(&retried);
420 let set_init = set_init.clone();
421
422 sendwrap_fn!(move || {
423 close();
424 explicitly_closed.store(false, std::sync::atomic::Ordering::Relaxed);
425 retried.store(0, std::sync::atomic::Ordering::Relaxed);
426 if init.get_value().is_none() && !url.get_untracked().is_empty() {
427 set_init(url.get_untracked());
428 }
429 if let Some(init) = init.get_value() {
430 init();
431 }
432 })
433 };
434
435 {
436 let close = close.clone();
437 let open = open.clone();
438 let set_init = set_init.clone();
439 Effect::watch(
440 move || url.get(),
441 move |url, prev_url, _| {
442 if url.is_empty() {
443 close();
444 } else if Some(url) != prev_url {
445 close();
446 set_init(url.to_owned());
447 open();
448 }
449 },
450 immediate,
451 );
452 }
453
454 on_cleanup(close.clone());
455 }
456
457 #[cfg(feature = "ssr")]
458 {
459 open = move || {};
460 close = move || {};
461
462 let _ = reconnect_limit;
463 let _ = reconnect_interval;
464 let _ = on_failed;
465 let _ = immediate;
466 let _ = named_events;
467 let _ = on_event;
468 let _ = with_credentials;
469
470 let _ = set_message;
471 let _ = set_ready_state;
472 let _ = set_error;
473 let _ = url;
474 }
475
476 UseEventSourceReturn {
477 message: message.into(),
478 ready_state: ready_state.into(),
479 error: error.into(),
480 open,
481 close,
482 }
483}
484
485/// Message received from the `EventSource` with transcoded data.
486#[derive(PartialEq)]
487pub struct UseEventSourceMessage<T, C>
488where
489 T: Clone + Send + Sync + 'static,
490 C: Decoder<T, Encoded = str> + Send + Sync,
491 C::Error: Send + Sync,
492{
493 pub event_type: String,
494 pub data: T,
495 pub last_event_id: String,
496 _marker: PhantomData<C>,
497}
498
499impl<T, C> Clone for UseEventSourceMessage<T, C>
500where
501 T: Clone + Send + Sync + 'static,
502 C: Decoder<T, Encoded = str> + Send + Sync,
503 C::Error: Send + Sync,
504{
505 fn clone(&self) -> Self {
506 Self {
507 event_type: self.event_type.clone(),
508 data: self.data.clone(),
509 last_event_id: self.last_event_id.clone(),
510 _marker: PhantomData,
511 }
512 }
513}
514
515impl<T, C> Debug for UseEventSourceMessage<T, C>
516where
517 T: Debug + Clone + Send + Sync + 'static,
518 C: Decoder<T, Encoded = str> + Send + Sync,
519 C::Error: Send + Sync,
520{
521 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
522 f.debug_struct("UseEventSourceMessage")
523 .field("data", &self.data)
524 .field("event_type", &self.event_type)
525 .field("last_event_id", &self.last_event_id)
526 .finish()
527 }
528}
529
530impl<T, C> TryFrom<&web_sys::MessageEvent> for UseEventSourceMessage<T, C>
531where
532 T: Clone + Send + Sync + 'static,
533 C: Decoder<T, Encoded = str> + Send + Sync,
534 C::Error: Send + Sync,
535{
536 type Error = UseEventSourceError<C::Error>;
537
538 fn try_from(message_event: &web_sys::MessageEvent) -> Result<Self, Self::Error> {
539 let data_string = message_event.data().as_string().unwrap_or_default();
540
541 let data = C::decode(&data_string).map_err(UseEventSourceError::Deserialize)?;
542
543 Ok(Self {
544 event_type: message_event.type_(),
545 data,
546 last_event_id: message_event.last_event_id(),
547 _marker: PhantomData,
548 })
549 }
550}
551
552impl<T, C> TryFrom<web_sys::Event> for UseEventSourceMessage<T, C>
553where
554 T: Clone + Send + Sync + 'static,
555 C: Decoder<T, Encoded = str> + Send + Sync,
556 C::Error: Send + Sync,
557{
558 type Error = UseEventSourceError<C::Error>;
559
560 fn try_from(event: web_sys::Event) -> Result<Self, Self::Error> {
561 let message_event = event
562 .dyn_into::<web_sys::MessageEvent>()
563 .map_err(|e| UseEventSourceError::CastToMessageEvent(e.type_()))?;
564
565 UseEventSourceMessage::try_from(&message_event)
566 }
567}
568
569/// Options for [`use_event_source_with_options`].
570#[derive(DefaultBuilder)]
571pub struct UseEventSourceOptions<T>
572where
573 T: 'static,
574{
575 /// Retry times. Defaults to `ReconnectLimit::Limited(3)`. Use `ReconnectLimit::Infinite` for
576 /// infinite retries.
577 reconnect_limit: ReconnectLimit,
578
579 /// Retry interval in ms. Defaults to 3000.
580 reconnect_interval: u64,
581
582 /// On maximum retry times reached.
583 on_failed: Arc<dyn Fn() + Send + Sync>,
584
585 /// If `true` the `EventSource` connection will immediately be opened when calling this function.
586 /// If `false` you have to manually call the `open` function.
587 /// Defaults to `true`.
588 immediate: bool,
589
590 /// List of named events to listen for on the `EventSource`.
591 #[builder(into)]
592 named_events: Vec<String>,
593
594 /// The `on_event` is called before processing any event inside of [`use_event_source`].
595 /// Return `UseEventSourceOnEventReturn::Ignore` to ignore further processing of the respective event
596 /// in [`use_event_source`], or `UseEventSourceOnEventReturn::Use` to process the event as usual.
597 ///
598 /// Beware that ignoring processing the `open` and `error` events may yield unexpected results.
599 ///
600 /// You may want to use [`UseEventSourceMessage::try_from()`] to access the event data.
601 ///
602 /// Default handler returns `UseEventSourceOnEventReturn::Use`.
603 on_event: Arc<dyn Fn(&web_sys::Event) -> UseEventSourceOnEventReturn + Send + Sync>,
604
605 /// If CORS should be set to `include` credentials. Defaults to `false`.
606 with_credentials: bool,
607
608 _marker: PhantomData<T>,
609}
610
611impl<T> Default for UseEventSourceOptions<T> {
612 fn default() -> Self {
613 Self {
614 reconnect_limit: ReconnectLimit::default(),
615 reconnect_interval: 3000,
616 on_failed: Arc::new(|| {}),
617 immediate: true,
618 named_events: vec![],
619 on_event: Arc::new(|_| UseEventSourceOnEventReturn::ProcessMessage),
620 with_credentials: false,
621 _marker: PhantomData,
622 }
623 }
624}
625
626/// Return type of the `on_event` handler in [`UseEventSourceOptions`].
627pub enum UseEventSourceOnEventReturn {
628 /// Ignore further processing of the message event in [`use_event_source`].
629 IgnoreProcessingMessage,
630 /// Use the default processing of the message event in [`use_event_source`].
631 ProcessMessage,
632}
633
634/// Return type of [`use_event_source`].
635pub struct UseEventSourceReturn<T, C, Err, OpenFn, CloseFn>
636where
637 T: Clone + Send + Sync + 'static,
638 C: Decoder<T, Encoded = str> + Send + Sync,
639 C::Error: Send + Sync,
640 Err: Send + Sync + 'static,
641 OpenFn: Fn() + Clone + Send + Sync + 'static,
642 CloseFn: Fn() + Clone + Send + Sync + 'static,
643{
644 /// The latest message
645 pub message: Signal<Option<UseEventSourceMessage<T, C>>>,
646
647 /// The current state of the connection,
648 pub ready_state: Signal<ConnectionReadyState>,
649
650 /// The current error
651 pub error: Signal<Option<UseEventSourceError<Err>>>,
652
653 /// (Re-)Opens the `EventSource` connection
654 /// If the current one is active, will close it before opening a new one.
655 pub open: OpenFn,
656
657 /// Closes the `EventSource` connection
658 pub close: CloseFn,
659}
660
661#[derive(Error, Debug)]
662pub enum UseEventSourceError<Err> {
663 #[error("Error event received")]
664 ErrorEvent,
665
666 #[error("Error decoding value")]
667 Deserialize(Err),
668
669 #[error("Error casting event '{0}' to MessageEvent")]
670 CastToMessageEvent(String),
671}