Skip to main content

iroh_http_core/endpoint/
lifecycle.rs

1//! Endpoint lifecycle: close, drain, serve-handle wiring, and events.
2//!
3//! Extracted from `mod.rs` to keep the façade ≤ 200 LoC (ADR-014 D1 AC #4).
4//! All methods here operate on [`SessionRuntime`] and [`Transport`] fields
5//! that live inside [`EndpointInner`].
6
7use super::IrohEndpoint;
8
9impl IrohEndpoint {
10    /// Graceful close: signal the serve loop to stop accepting, wait for
11    /// in-flight requests to drain, then close the QUIC endpoint.
12    pub async fn close(&self) {
13        let handle = self
14            .inner
15            .session
16            .serve_handle
17            .lock()
18            .unwrap_or_else(|e| e.into_inner())
19            .take();
20        if let Some(h) = handle {
21            h.drain().await;
22        }
23        self.inner.transport.ep.close().await;
24        let _ = self.inner.session.closed_tx.send(true);
25    }
26
27    /// Immediate close: abort the serve loop with no drain period.
28    pub async fn close_force(&self) {
29        let handle = self
30            .inner
31            .session
32            .serve_handle
33            .lock()
34            .unwrap_or_else(|e| e.into_inner())
35            .take();
36        if let Some(h) = handle {
37            h.abort();
38        }
39        self.inner.transport.ep.close().await;
40        let _ = self.inner.session.closed_tx.send(true);
41    }
42
43    /// Wait until this endpoint has been closed. Returns immediately if already closed.
44    pub async fn wait_closed(&self) {
45        let mut rx = self.inner.session.closed_rx.clone();
46        let _ = rx.wait_for(|v| *v).await;
47    }
48
49    /// Store a serve handle so that `close()` can drain it.
50    ///
51    /// If `stop_serve()` was called before this (the JS abort raced ahead of
52    /// the napi async setup), the handle is immediately shut down so the
53    /// serve loop drains without a second `stop_serve()` call.
54    pub fn set_serve_handle(&self, handle: crate::http::server::ServeHandle) {
55        // Clear the early-stop flag and check if it was set.
56        let was_stopped = self
57            .inner
58            .session
59            .serve_stopped_early
60            .swap(false, std::sync::atomic::Ordering::AcqRel);
61
62        *self
63            .inner
64            .session
65            .serve_done_rx
66            .lock()
67            .unwrap_or_else(|e| e.into_inner()) = Some(handle.subscribe_done());
68
69        if was_stopped {
70            handle.shutdown();
71        }
72
73        *self
74            .inner
75            .session
76            .serve_handle
77            .lock()
78            .unwrap_or_else(|e| e.into_inner()) = Some(handle);
79    }
80
81    /// Signal the serve loop to stop accepting new connections.
82    pub fn stop_serve(&self) {
83        let guard = self
84            .inner
85            .session
86            .serve_handle
87            .lock()
88            .unwrap_or_else(|e| e.into_inner());
89        if let Some(h) = guard.as_ref() {
90            h.shutdown();
91        } else {
92            // Handle not registered yet — set a flag so that
93            // `set_serve_handle` will shut down as soon as it arrives.
94            self.inner
95                .session
96                .serve_stopped_early
97                .store(true, std::sync::atomic::Ordering::Release);
98        }
99    }
100
101    /// Wait until the serve loop has fully exited.
102    pub async fn wait_serve_stop(&self) {
103        let rx = self
104            .inner
105            .session
106            .serve_done_rx
107            .lock()
108            .unwrap_or_else(|e| e.into_inner())
109            .clone();
110        if let Some(mut rx) = rx {
111            let _ = rx.wait_for(|v| *v).await;
112        }
113    }
114
115    /// Take the transport event receiver, handing it off to a platform drain task.
116    /// May only be called once per endpoint.
117    pub fn subscribe_events(
118        &self,
119    ) -> Option<tokio::sync::mpsc::Receiver<crate::http::events::TransportEvent>> {
120        self.inner
121            .session
122            .event_rx
123            .lock()
124            .unwrap_or_else(|e| e.into_inner())
125            .take()
126    }
127}