MCPcopy Create free account
hub / github.com/clockworklabs/SpacetimeDB / ws_main_loop

Function ws_main_loop

crates/client-api/src/routes/subscribe.rs:613–709  ·  view source on GitHub ↗

The main `select!` loop of the websocket client actor. > This function is defined standalone with generic parameters so that its > behavior can be tested in isolation, not requiring I/O and allowing to > mock effects easily. The loop's responsibilities are: - Drive the tasks handling the send and receive ends of the websockets to completion, terminating when either of them completes. - Termina

(
    state: Arc<ActorState>,
    hotswap: impl Fn() -> HotswapWatcher,
    idle_timer: impl Future<Output = ()>,
    mut send_task: JoinHandle<()>,
    mut recv_task: JoinHandle<()>,
    unordered_tx

Source from the content-addressed store, hash-verified

611///
612/// [close handshake]: https://datatracker.ietf.org/doc/html/rfc6455#section-7
613async fn ws_main_loop<HotswapWatcher>(
614 state: Arc<ActorState>,
615 hotswap: impl Fn() -> HotswapWatcher,
616 idle_timer: impl Future<Output = ()>,
617 mut send_task: JoinHandle<()>,
618 mut recv_task: JoinHandle<()>,
619 unordered_tx: impl Fn(UnorderedWsMessage),
620) where
621 HotswapWatcher: Future<Output = Result<(), NoSuchModule>>,
622{
623 // Ensure we terminate both tasks if either exits.
624 let abort_send = send_task.abort_handle();
625 let abort_recv = recv_task.abort_handle();
626 defer! {
627 abort_send.abort();
628 abort_recv.abort();
629 };
630 // Set up the ping interval.
631 let mut ping_interval = tokio::time::interval(state.config.ping_interval);
632 // Arm the first hotswap watcher.
633 let watch_hotswap = hotswap();
634
635 pin_mut!(watch_hotswap);
636 pin_mut!(idle_timer);
637
638 loop {
639 let closed = state.closed();
640
641 tokio::select! {
642 // Drive send and receive tasks to completion,
643 // propagating panics.
644 //
645 // If either task completes,
646 // the connection is considered closed and we break the loop.
647 //
648 // NOTE: We don't abort the tasks until this function returns,
649 // so the `Err` can't contain an `is_cancelled()` value.
650 //
651 // Even if the tasks were cancelled (e.g. if the caller retains
652 // [`tokio::task::AbortHandle`]s), the reasonable thing to do is to
653 // exit the loop as if the tasks completed normally.
654 res = &mut send_task => {
655 if let Err(e) = res
656 && e.is_panic() {
657 panic::resume_unwind(e.into_panic())
658 }
659 break;
660 },
661 res = &mut recv_task => {
662 if let Err(e) = res
663 && e.is_panic() {
664 panic::resume_unwind(e.into_panic())
665 }
666 break;
667 },
668
669 // Exit if we haven't heard from the client for too long.
670 _ = &mut idle_timer => {

Calls 3

intervalFunction · 0.85
abort_handleMethod · 0.45
closedMethod · 0.45

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…