1pub mod whep;
10pub mod whip;
11
12mod mux;
13
14use std::collections::HashMap;
15use std::net::SocketAddr;
16use std::sync::{Arc, Mutex};
17
18use axum::Router;
19use axum::extract::{Path, State};
20use axum::http::{HeaderValue, StatusCode, Uri};
21use tokio::sync::{OnceCell, oneshot};
22
23use crate::{Error, Result};
24use mux::Mux;
25
26pub struct Response {
30 pub resource_id: String,
32 pub answer: String,
34 session: AcceptedSession,
35}
36
37impl Response {
38 pub(crate) fn new(
40 server: Server,
41 resource_id: String,
42 answer: String,
43 session: crate::session::Session,
44 registration: mux::Registration,
45 cancel: oneshot::Receiver<()>,
46 role: &'static str,
47 ) -> Self {
48 Self {
49 resource_id: resource_id.clone(),
50 answer,
51 session: AcceptedSession {
52 server,
53 resource_id,
54 session: Some(session),
55 registration: Some(registration),
56 cancel: Some(cancel),
57 role,
58 },
59 }
60 }
61
62 pub async fn run(self) -> Result<()> {
64 self.session.run().await
65 }
66}
67
68impl std::fmt::Debug for Response {
69 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
70 f.debug_struct("Response")
71 .field("resource_id", &self.resource_id)
72 .field("answer", &self.answer)
73 .finish_non_exhaustive()
74 }
75}
76
77struct AcceptedSession {
78 server: Server,
79 resource_id: String,
80 session: Option<crate::session::Session>,
81 registration: Option<mux::Registration>,
82 cancel: Option<oneshot::Receiver<()>>,
83 role: &'static str,
84}
85
86impl AcceptedSession {
87 async fn run(mut self) -> Result<()> {
88 let session = self.session.take().expect("accepted session missing driver");
89 let registration = self
90 .registration
91 .take()
92 .expect("accepted session missing mux registration");
93 let cancel = self.cancel.take().expect("accepted session missing cancel receiver");
94
95 let result = {
96 let _registration = registration;
97 tokio::select! {
98 res = session.run() => {
99 crate::session::log_session_end(self.role, &res);
100 res
101 }
102 _ = cancel => {
103 tracing::debug!(role = self.role, "webrtc session terminated by DELETE");
104 Ok(())
105 }
106 }
107 };
108 normalize_session_result(result)
109 }
110}
111
112impl Drop for AcceptedSession {
113 fn drop(&mut self) {
114 self.server.unregister_session(&self.resource_id);
115 }
116}
117
118fn normalize_session_result(result: Result<()>) -> Result<()> {
119 match result {
120 Ok(()) | Err(Error::SessionClosed) => Ok(()),
121 Err(err) => Err(err),
122 }
123}
124
125pub(crate) fn session_location(uri: &Uri, resource_id: &str) -> Option<HeaderValue> {
126 let base = uri.path().trim_end_matches('/');
127 let path = if base.is_empty() {
128 format!("/{resource_id}")
129 } else {
130 format!("{base}/{resource_id}")
131 };
132 HeaderValue::from_str(&path).ok()
133}
134
135#[derive(Clone, Debug)]
137#[non_exhaustive]
138pub struct Config {
139 pub ice_candidates: Vec<SocketAddr>,
147
148 pub udp_bind: SocketAddr,
153}
154
155impl Default for Config {
156 fn default() -> Self {
157 Self {
158 ice_candidates: Vec::new(),
159 udp_bind: SocketAddr::from(([0, 0, 0, 0], 0)),
160 }
161 }
162}
163
164#[derive(Clone)]
170pub struct Server {
171 inner: Arc<Inner>,
172}
173
174struct Inner {
175 config: Config,
176 publisher: moq_net::OriginProducer,
177 subscriber: moq_net::OriginConsumer,
179 mux: OnceCell<Mux>,
182 sessions: Mutex<HashMap<String, oneshot::Sender<()>>>,
186}
187
188impl Server {
189 pub fn new(config: Config, publisher: moq_net::OriginProducer, subscriber: moq_net::OriginConsumer) -> Self {
192 Self {
193 inner: Arc::new(Inner {
194 config,
195 publisher,
196 subscriber,
197 mux: OnceCell::new(),
198 sessions: Mutex::new(HashMap::new()),
199 }),
200 }
201 }
202
203 pub(crate) async fn mux(&self) -> Result<&Mux> {
205 self.inner
206 .mux
207 .get_or_try_init(|| Mux::bind(self.inner.config.udp_bind, &self.inner.config.ice_candidates))
208 .await
209 }
210
211 pub fn publish_router(&self) -> Router {
219 whip::router(self.clone())
220 }
221
222 pub fn subscribe_router(&self) -> Router {
230 whep::router(self.clone())
231 }
232
233 pub(crate) fn publisher(&self) -> &moq_net::OriginProducer {
234 &self.inner.publisher
235 }
236
237 pub(crate) fn subscriber(&self) -> &moq_net::OriginConsumer {
238 &self.inner.subscriber
239 }
240
241 pub(crate) fn register_session(&self, resource_id: String) -> oneshot::Receiver<()> {
245 let (tx, rx) = oneshot::channel();
246 self.inner.sessions.lock().unwrap().insert(resource_id, tx);
247 rx
248 }
249
250 pub(crate) fn unregister_session(&self, resource_id: &str) {
252 self.inner.sessions.lock().unwrap().remove(resource_id);
253 }
254
255 pub fn terminate(&self, resource_id: &str) -> bool {
262 if let Some(cancel) = self.inner.sessions.lock().unwrap().remove(resource_id) {
263 let _ = cancel.send(());
264 true
265 } else {
266 false
267 }
268 }
269}
270
271pub(crate) async fn delete(State(server): State<Server>, Path(path): Path<String>) -> StatusCode {
274 match crate::sdp::parse_resource_id(&path) {
275 Ok(id) if server.terminate(&id.to_string()) => StatusCode::OK,
276 Ok(_) => StatusCode::NOT_FOUND,
277 Err(_) => StatusCode::BAD_REQUEST,
278 }
279}
280
281#[cfg(test)]
282mod tests {
283 use super::*;
284
285 fn server() -> Server {
286 let publisher = moq_net::Origin::random().produce();
287 let subscriber = moq_net::Origin::random().produce().consume();
288 Server::new(Config::default(), publisher, subscriber)
289 }
290
291 #[test]
292 fn terminate_unknown_session_is_false() {
293 assert!(!server().terminate("00000000-0000-0000-0000-000000000000"));
294 }
295
296 #[test]
297 fn terminate_registered_session_once() {
298 let server = server();
299 let id = "11111111-1111-1111-1111-111111111111";
300 let _cancel = server.register_session(id.to_string());
301 assert!(server.terminate(id), "first terminate finds the session");
302 assert!(!server.terminate(id), "second terminate is a no-op");
303 }
304
305 #[test]
306 fn unregister_drops_the_entry() {
307 let server = server();
308 let id = "22222222-2222-2222-2222-222222222222";
309 let _cancel = server.register_session(id.to_string());
310 server.unregister_session(id);
311 assert!(!server.terminate(id), "unregistered session can't be terminated");
312 }
313
314 #[test]
315 fn peer_close_is_a_successful_session_result() {
316 assert!(normalize_session_result(Err(Error::SessionClosed)).is_ok());
317 }
318
319 #[test]
320 fn session_location_preserves_mount_path() {
321 let uri: Uri = "/whip/live/cam0?token=secret".parse().unwrap();
322 let location = session_location(&uri, "session-id").expect("header value");
323 assert_eq!(location, "/whip/live/cam0/session-id");
324 }
325}