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}