1use std::net;
2#[cfg(any(test, all(feature = "uds", unix)))]
3use std::path::PathBuf;
4
5#[cfg(feature = "iroh")]
6use crate::iroh;
7use crate::{Error, QuicBackend};
8use moq_net::Session;
9use url::Url;
10
11#[cfg(any(
14 feature = "noq",
15 feature = "quinn",
16 feature = "quiche",
17 feature = "iroh",
18 feature = "websocket"
19))]
20use futures::FutureExt;
21use futures::future::BoxFuture;
22use futures::stream::FuturesUnordered;
23use futures::stream::StreamExt;
24
25#[derive(clap::Args, Clone, Debug, Default, serde::Serialize, serde::Deserialize)]
27#[serde(deny_unknown_fields, default)]
28#[non_exhaustive]
29pub struct ServerConfig {
30 #[serde(alias = "listen")]
38 #[arg(id = "server-bind", long = "server-bind", alias = "listen", env = "MOQ_SERVER_BIND")]
39 pub bind: Option<String>,
40
41 #[cfg(feature = "tcp")]
44 #[command(flatten)]
45 #[serde(default)]
46 pub tcp: crate::tcp::Config,
47
48 #[cfg(all(feature = "uds", unix))]
51 #[command(flatten)]
52 #[serde(default)]
53 pub unix: crate::unix::Config,
54
55 #[arg(id = "server-backend", long = "server-backend", env = "MOQ_SERVER_BACKEND")]
58 pub backend: Option<QuicBackend>,
59
60 #[command(flatten)]
63 #[serde(default)]
64 pub quic: crate::quic::Server,
65
66 #[serde(default, skip_serializing_if = "Vec::is_empty")]
72 #[arg(
73 id = "server-version",
74 long = "server-version",
75 env = "MOQ_SERVER_VERSION",
76 value_parser = crate::version_parser(),
77 )]
78 pub version: Vec<moq_net::Version>,
79
80 #[command(flatten)]
83 #[serde(default)]
84 pub tls: crate::tls::Server,
85}
86
87impl ServerConfig {
88 pub fn init(self) -> crate::Result<Server> {
90 Server::new(self)
91 }
92
93 pub fn versions(&self) -> moq_net::Versions {
95 if self.version.is_empty() {
96 moq_net::Versions::all()
97 } else {
98 moq_net::Versions::from(self.version.clone())
99 }
100 }
101
102 #[allow(unused_mut)]
107 fn has_stream_listener(&self) -> bool {
108 let mut has = false;
109 #[cfg(feature = "tcp")]
110 {
111 has |= self.tcp.bind.is_some();
112 }
113 #[cfg(all(feature = "uds", unix))]
114 {
115 has |= self.unix.bind.is_some();
116 }
117 has
118 }
119}
120
121#[cfg(any(feature = "noq", feature = "quinn", feature = "quiche"))]
123pub(crate) const DEFAULT_BIND: &str = "[::]:443";
124
125pub struct Server {
131 moq: moq_net::Server,
132 versions: moq_net::Versions,
133 accept: FuturesUnordered<BoxFuture<'static, crate::Result<Request>>>,
134 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
135 streams: StreamListeners,
136 #[cfg(feature = "iroh")]
137 iroh: Option<iroh::Endpoint>,
138 #[cfg(feature = "noq")]
139 noq: Option<crate::noq::NoqServer>,
140 #[cfg(feature = "quinn")]
141 quinn: Option<crate::quinn::QuinnServer>,
142 #[cfg(feature = "quiche")]
143 quiche: Option<crate::quiche::QuicheServer>,
144 #[cfg(feature = "websocket")]
145 websocket: Option<crate::websocket::Listener>,
146}
147
148impl Server {
149 pub fn new(config: ServerConfig) -> crate::Result<Self> {
154 #[cfg(any(feature = "noq", feature = "quinn", feature = "quiche"))]
157 let backend = config.backend.clone().unwrap_or_else(crate::default_quic_backend);
158
159 let versions = config.versions();
160
161 config.quic.validate()?;
165
166 let build_quic = config.bind.is_some() || !config.has_stream_listener();
167 #[cfg(not(any(feature = "noq", feature = "quinn", feature = "quiche")))]
168 if config.bind.is_some() {
169 return Err(Error::NoBackend(
170 "--server-bind requires a noq, quinn, or quiche backend feature",
171 ));
172 }
173
174 if build_quic && !config.tls.root.is_empty() {
175 #[cfg(any(feature = "noq", feature = "quinn", feature = "quiche"))]
178 let mtls_supported = match backend {
179 #[cfg(feature = "quinn")]
180 QuicBackend::Quinn => true,
181 #[cfg(feature = "noq")]
182 QuicBackend::Noq => true,
183 #[cfg(feature = "quiche")]
184 QuicBackend::Quiche => true,
185 #[allow(unreachable_patterns)]
186 _ => false,
187 };
188 #[cfg(not(any(feature = "noq", feature = "quinn", feature = "quiche")))]
189 let mtls_supported = false;
190
191 if !mtls_supported {
192 return Err(Error::MtlsUnsupported);
193 }
194 }
195
196 #[cfg(feature = "noq")]
197 #[allow(unreachable_patterns)]
198 let noq = match backend {
199 QuicBackend::Noq if build_quic => Some(crate::noq::NoqServer::new(config.clone())?),
200 _ => None,
201 };
202
203 #[cfg(feature = "quinn")]
204 #[allow(unreachable_patterns)]
205 let quinn = match backend {
206 QuicBackend::Quinn if build_quic => Some(crate::quinn::QuinnServer::new(config.clone())?),
207 _ => None,
208 };
209
210 #[cfg(feature = "quiche")]
211 let quiche = match backend {
212 QuicBackend::Quiche if build_quic => Some(crate::quiche::QuicheServer::new(config.clone())?),
213 _ => None,
214 };
215
216 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
218 let mut stream_binds = Vec::new();
219 #[cfg(feature = "tcp")]
220 if let Some(addr) = config.tcp.bind {
221 stream_binds.push(StreamBind::Tcp(addr));
222 }
223 #[cfg(all(feature = "uds", unix))]
224 if let Some(path) = config.unix.bind.clone() {
225 stream_binds.push(StreamBind::Unix(path));
226 }
227 #[cfg(all(feature = "uds", unix))]
229 let unix_allow = config.unix.allow.clone().filter(|allow| !allow.is_empty());
230 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
231 let streams = StreamListeners::new(
232 stream_binds,
233 stream_versions(&versions),
234 #[cfg(all(feature = "uds", unix))]
235 unix_allow,
236 );
237
238 Ok(Server {
239 accept: Default::default(),
240 moq: moq_net::Server::new().with_versions(versions.clone()),
241 versions,
242 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
243 streams,
244 #[cfg(feature = "iroh")]
245 iroh: None,
246 #[cfg(feature = "noq")]
247 noq,
248 #[cfg(feature = "quinn")]
249 quinn,
250 #[cfg(feature = "quiche")]
251 quiche,
252 #[cfg(feature = "websocket")]
253 websocket: None,
254 })
255 }
256
257 #[cfg(feature = "websocket")]
263 pub fn with_websocket(mut self, websocket: crate::websocket::Listener) -> Self {
264 self.websocket = Some(websocket);
265 self
266 }
267
268 #[cfg(feature = "iroh")]
270 pub fn with_iroh(mut self, iroh: iroh::Endpoint) -> Self {
271 self.iroh = Some(iroh);
272 self
273 }
274
275 pub fn with_publisher(mut self, publish: impl moq_net::Consume<moq_net::origin::Consumer>) -> Self {
277 self.moq = self.moq.with_publisher(publish);
278 self
279 }
280
281 pub fn with_subscriber(mut self, subscribe: moq_net::origin::Producer) -> Self {
283 self.moq = self.moq.with_subscriber(subscribe);
284 self
285 }
286
287 pub fn with_stats(mut self, stats: moq_net::stats::Session) -> Self {
290 self.moq = self.moq.with_stats(stats);
291 self
292 }
293
294 pub async fn serve_publish(self, origin: moq_net::origin::Consumer) -> crate::Result<()> {
301 self.with_publisher(origin).serve().await
302 }
303
304 pub async fn serve_consume(self, origin: moq_net::origin::Producer) -> crate::Result<()> {
308 self.with_subscriber(origin).serve().await
309 }
310
311 pub async fn serve_both(
318 self,
319 publish: moq_net::origin::Consumer,
320 subscribe: moq_net::origin::Producer,
321 ) -> crate::Result<()> {
322 self.with_publisher(publish).with_subscriber(subscribe).serve().await
323 }
324
325 async fn serve(mut self) -> crate::Result<()> {
329 if let Ok(addr) = self.local_addr() {
330 tracing::info!(%addr, "listening");
331 }
332 while let Some(request) = self.accept().await {
333 tokio::spawn(async move {
334 if let Err(err) = serve_session(request).await {
335 tracing::warn!(%err, "session ended with error");
336 }
337 });
338 }
339 Ok(())
340 }
341
342 pub fn certificates(&self) -> crate::tls::Certificates {
351 #[cfg(feature = "noq")]
352 if let Some(noq) = self.noq.as_ref() {
353 return noq.certificates();
354 }
355 #[cfg(feature = "quinn")]
356 if let Some(quinn) = self.quinn.as_ref() {
357 return quinn.certificates();
358 }
359 #[cfg(feature = "quiche")]
360 if let Some(quiche) = self.quiche.as_ref() {
361 return quiche.certificates();
362 }
363 crate::tls::Certificates::empty()
365 }
366
367 #[cfg(not(any(
368 feature = "noq",
369 feature = "quinn",
370 feature = "quiche",
371 feature = "iroh",
372 feature = "websocket",
373 feature = "tcp",
374 all(feature = "uds", unix)
375 )))]
376 pub async fn accept(&mut self) -> Option<Request> {
380 unreachable!("no transport compiled; enable a QUIC backend, websocket, tcp, or uds feature");
381 }
382
383 pub fn accept_health(&self) -> Vec<crate::accept::Health> {
395 #[allow(unused_mut)]
396 let mut health = Vec::new();
397 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
398 health.extend(self.streams.health.iter().cloned());
399 #[cfg(feature = "websocket")]
400 health.extend(self.websocket.as_ref().map(|ws| ws.accept_health()));
401 health
402 }
403
404 pub async fn listen(&mut self) -> crate::Result<()> {
421 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
422 self.streams.ensure_started(self.moq.clone()).await?;
423 Ok(())
424 }
425
426 #[cfg(any(
437 feature = "noq",
438 feature = "quinn",
439 feature = "quiche",
440 feature = "iroh",
441 feature = "websocket",
442 feature = "tcp",
443 all(feature = "uds", unix)
444 ))]
445 pub async fn accept(&mut self) -> Option<Request> {
446 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
450 if let Err(err) = self.streams.ensure_started(self.moq.clone()).await {
451 tracing::error!(%err, "failed to bind stream listener");
452 return None;
453 }
454
455 loop {
456 #[cfg(feature = "noq")]
458 let noq_accept = async {
459 #[cfg(feature = "noq")]
460 if let Some(noq) = self.noq.as_mut() {
461 return noq.accept().await;
462 }
463 None
464 };
465 #[cfg(not(feature = "noq"))]
466 let noq_accept = async { None::<()> };
467
468 #[cfg(feature = "iroh")]
469 let iroh_accept = async {
470 #[cfg(feature = "iroh")]
471 if let Some(endpoint) = self.iroh.as_mut() {
472 return endpoint.accept().await;
473 }
474 None
475 };
476 #[cfg(not(feature = "iroh"))]
477 let iroh_accept = async { None::<()> };
478
479 #[cfg(feature = "quinn")]
480 let quinn_accept = async {
481 #[cfg(feature = "quinn")]
482 if let Some(quinn) = self.quinn.as_mut() {
483 return quinn.accept().await;
484 }
485 None
486 };
487 #[cfg(not(feature = "quinn"))]
488 let quinn_accept = async { None::<()> };
489
490 #[cfg(feature = "quiche")]
491 let quiche_accept = async {
492 #[cfg(feature = "quiche")]
493 if let Some(quiche) = self.quiche.as_mut() {
494 return quiche.accept().await;
495 }
496 None
497 };
498 #[cfg(not(feature = "quiche"))]
499 let quiche_accept = async { None::<()> };
500
501 #[cfg(feature = "websocket")]
502 let ws_ref = self.websocket.as_ref();
503 #[cfg(feature = "websocket")]
504 let ws_accept = async {
505 match ws_ref {
506 Some(ws) => ws.accept_with_url().await,
507 None => std::future::pending().await,
508 }
509 };
510 #[cfg(not(feature = "websocket"))]
511 let ws_accept = std::future::pending::<Option<crate::Result<()>>>();
512
513 #[allow(unused_variables)]
514 let server = self.moq.clone();
515 #[allow(unused_variables)]
516 let versions = self.versions.clone();
517
518 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
520 let stream_accept = self.streams.recv();
521 #[cfg(not(any(feature = "tcp", all(feature = "uds", unix))))]
522 let stream_accept = std::future::pending::<Option<Request>>();
523
524 tokio::select! {
525 Some(request) = stream_accept => {
526 return Some(request);
527 }
528 Some(_conn) = noq_accept => {
529 #[cfg(feature = "noq")]
530 {
531 let alpns = versions.alpns();
532 self.accept.push(async move {
533 let (session, url, identity) = super::noq::accept(_conn, alpns).await?;
537 let request = server.accept_request(session).await?;
538 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Noq(Box::new(request)) })
539 }.boxed());
540 }
541 }
542 Some(_conn) = quinn_accept => {
543 #[cfg(feature = "quinn")]
544 {
545 let alpns = versions.alpns();
546 self.accept.push(async move {
547 let (session, url, identity) = super::quinn::accept(_conn, alpns).await?;
548 let request = server.accept_request(session).await?;
549 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Quinn(Box::new(request)) })
550 }.boxed());
551 }
552 }
553 Some(_conn) = quiche_accept => {
554 #[cfg(feature = "quiche")]
555 {
556 let alpns = versions.alpns();
557 self.accept.push(async move {
558 let (session, url, identity) = super::quiche::accept(_conn, alpns).await?;
559 let request = server.accept_request(session).await?;
560 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Quiche(Box::new(request)) })
561 }.boxed());
562 }
563 }
564 Some(_conn) = iroh_accept => {
565 #[cfg(feature = "iroh")]
566 self.accept.push(async move {
567 let (session, url, identity) = super::iroh::accept(_conn).await?;
568 let request = server.accept_request(session).await?;
569 Ok(Request { transport: Transport::Iroh, url, identity, kind: RequestKind::Iroh(Box::new(request)) })
570 }.boxed());
571 }
572 Some(_res) = ws_accept => {
573 #[cfg(feature = "websocket")]
574 match _res {
575 Ok((session, url)) => {
576 self.accept.push(async move {
579 let request = server.accept_request(session).await?;
580 Ok(Request { transport: Transport::WebSocket, url: Some(url), identity: None, kind: RequestKind::Qmux(Box::new(request)) })
581 }.boxed());
582 }
583 Err(err) => tracing::debug!(%err, "WebSocket upgrade failed"),
587 }
588 }
589 Some(res) = self.accept.next() => {
590 match res {
591 Ok(session) => return Some(session),
592 Err(err) => tracing::debug!(%err, "failed to accept session"),
593 }
594 }
595 _ = tokio::signal::ctrl_c() => {
596 self.close().await;
597 return None;
598 }
599 }
600 }
601 }
602
603 #[cfg(feature = "iroh")]
605 pub fn iroh_endpoint(&self) -> Option<&iroh::Endpoint> {
606 self.iroh.as_ref()
607 }
608
609 pub fn local_addr(&self) -> crate::Result<net::SocketAddr> {
615 #[cfg(feature = "noq")]
616 if let Some(noq) = self.noq.as_ref() {
617 return Ok(noq.local_addr()?);
618 }
619 #[cfg(feature = "quinn")]
620 if let Some(quinn) = self.quinn.as_ref() {
621 return Ok(quinn.local_addr()?);
622 }
623 #[cfg(feature = "quiche")]
624 if let Some(quiche) = self.quiche.as_ref() {
625 return Ok(quiche.local_addr()?);
626 }
627 Err(Error::NoBackend("no QUIC listener configured"))
629 }
630
631 #[cfg(feature = "websocket")]
634 pub fn websocket_local_addr(&self) -> Option<net::SocketAddr> {
635 self.websocket.as_ref().and_then(|ws| ws.local_addr().ok())
636 }
637
638 pub async fn close(&mut self) {
643 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
644 self.streams.close().await;
645 #[cfg(feature = "noq")]
646 if let Some(noq) = self.noq.as_mut() {
647 noq.close();
648 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
649 }
650 #[cfg(feature = "quinn")]
651 if let Some(quinn) = self.quinn.as_mut() {
652 quinn.close();
653 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
654 }
655 #[cfg(feature = "quiche")]
656 if let Some(quiche) = self.quiche.as_mut() {
657 quiche.close();
658 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
659 }
660 #[cfg(feature = "iroh")]
661 if let Some(iroh) = self.iroh.take() {
662 iroh.close().await;
663 }
664 #[cfg(feature = "websocket")]
665 {
666 let _ = self.websocket.take();
667 }
668 }
669}
670
671async fn serve_session(request: Request) -> crate::Result<()> {
673 let session = request.ok().await?;
674 Err(session.closed().await.into())
675}
676
677#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
683fn stream_versions(base: &moq_net::Versions) -> moq_net::Versions {
684 let mut versions: Vec<moq_net::Version> = base.iter().copied().collect();
685 if let Ok(lite05) = "moq-lite-05".parse::<moq_net::Version>()
686 && !versions.contains(&lite05)
687 {
688 versions.push(lite05);
689 }
690 moq_net::Versions::from(versions)
691}
692
693#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
695#[derive(Clone)]
696enum StreamBind {
697 #[cfg(feature = "tcp")]
698 Tcp(net::SocketAddr),
699 #[cfg(all(feature = "uds", unix))]
700 Unix(PathBuf),
701}
702
703#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
704impl StreamBind {
705 fn name(&self) -> &'static str {
707 match self {
708 #[cfg(feature = "tcp")]
709 Self::Tcp(_) => "tcp",
710 #[cfg(all(feature = "uds", unix))]
711 Self::Unix(_) => "unix",
712 }
713 }
714}
715
716#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
723struct StreamListeners {
724 binds: Vec<StreamBind>,
725 health: Vec<crate::accept::Health>,
729 versions: moq_net::Versions,
730 #[cfg(all(feature = "uds", unix))]
731 unix_allow: Option<crate::unix::Allow>,
732 rx: Option<tokio::sync::mpsc::Receiver<Request>>,
733 tasks: Vec<tokio::task::JoinHandle<()>>,
734}
735
736#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
737impl StreamListeners {
738 fn new(
739 binds: Vec<StreamBind>,
740 versions: moq_net::Versions,
741 #[cfg(all(feature = "uds", unix))] unix_allow: Option<crate::unix::Allow>,
742 ) -> Self {
743 let health = binds
744 .iter()
745 .map(|bind| crate::accept::Health::new(bind.name()))
746 .collect();
747 Self {
748 binds,
749 health,
750 versions,
751 #[cfg(all(feature = "uds", unix))]
752 unix_allow,
753 rx: None,
754 tasks: Vec::new(),
755 }
756 }
757
758 async fn ensure_started(&mut self, server: moq_net::Server) -> crate::Result<()> {
763 if self.rx.is_some() || self.binds.is_empty() {
764 return Ok(());
765 }
766
767 let server = server.with_versions(self.versions.clone());
770
771 let (tx, rx) = tokio::sync::mpsc::channel(16);
772 if let Err(err) = self.start(&server, &tx).await {
773 for task in self.tasks.drain(..) {
779 task.abort();
780 }
781 return Err(err);
782 }
783
784 self.rx = Some(rx);
785 Ok(())
786 }
787
788 async fn start(&mut self, server: &moq_net::Server, tx: &tokio::sync::mpsc::Sender<Request>) -> crate::Result<()> {
790 let binds = self.binds.clone();
793 let health = self.health.clone();
794 for (bind, health) in binds.into_iter().zip(health) {
795 let alpns = self.versions.alpns();
796 match bind {
797 #[cfg(feature = "tcp")]
798 StreamBind::Tcp(addr) => {
799 if !addr.ip().is_loopback() {
800 tracing::warn!(%addr, "tcp listener bound to a non-loopback address; qmux is UNENCRYPTED, ensure the network is trusted");
801 }
802 let listener = crate::tcp::Listener::bind(addr)
803 .await?
804 .with_protocols(alpns)
805 .with_accept_health(health);
806 tracing::info!(%addr, "listening (tcp)");
807 self.tasks.push(spawn_tcp_loop(listener, server.clone(), tx.clone()));
808 }
809 #[cfg(all(feature = "uds", unix))]
810 StreamBind::Unix(path) => {
811 let listener = crate::unix::Listener::bind(&path)
812 .await?
813 .with_protocols(alpns)
814 .with_accept_health(health);
815 listener.set_mode(0o666)?;
818 tracing::info!(path = %path.display(), allow = ?self.unix_allow, "listening (unix)");
819 self.tasks.push(spawn_unix_loop(
820 listener,
821 server.clone(),
822 self.unix_allow.clone(),
823 tx.clone(),
824 ));
825 }
826 }
827 }
828
829 Ok(())
830 }
831
832 async fn recv(&mut self) -> Option<Request> {
834 match self.rx.as_mut() {
835 Some(rx) => rx.recv().await,
836 None => std::future::pending().await,
837 }
838 }
839
840 async fn close(&mut self) {
842 self.binds.clear();
843 self.rx = None;
844 let tasks = std::mem::take(&mut self.tasks);
845 for task in tasks {
846 task.abort();
847 let _ = task.await;
848 }
849 }
850}
851
852#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
853impl Drop for StreamListeners {
854 fn drop(&mut self) {
855 for task in &self.tasks {
857 task.abort();
858 }
859 }
860}
861
862#[cfg(feature = "tcp")]
863fn spawn_tcp_loop(
864 listener: crate::tcp::Listener,
865 server: moq_net::Server,
866 tx: tokio::sync::mpsc::Sender<Request>,
867) -> tokio::task::JoinHandle<()> {
868 tokio::spawn(async move {
869 loop {
870 match listener.accept().await {
871 Some(Ok(session)) => spawn_stream_request(session, Transport::Tcp, server.clone(), tx.clone()),
872 Some(Err(err)) => tracing::warn!(%err, "tcp qmux handshake failed"),
875 None => break,
876 }
877 }
878 })
879}
880
881#[cfg(all(feature = "uds", unix))]
882fn spawn_unix_loop(
883 listener: crate::unix::Listener,
884 server: moq_net::Server,
885 allow: Option<crate::unix::Allow>,
886 tx: tokio::sync::mpsc::Sender<Request>,
887) -> tokio::task::JoinHandle<()> {
888 tokio::spawn(async move {
889 loop {
890 match listener.accept().await {
891 Some(Ok((session, cred))) => {
892 if let Some(allow) = &allow
894 && !allow.permits(&cred)
895 {
896 tracing::warn!(uid = cred.uid, gid = cred.gid, pid = ?cred.pid, "unix connection rejected by allow list");
897 continue;
898 }
899 spawn_stream_request(session, Transport::Unix, server.clone(), tx.clone());
900 }
901 Some(Err(err)) => tracing::warn!(%err, "unix qmux handshake failed"),
903 None => break,
904 }
905 }
906 })
907}
908
909#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
912fn spawn_stream_request(
913 session: qmux::Session,
914 transport: Transport,
915 server: moq_net::Server,
916 tx: tokio::sync::mpsc::Sender<Request>,
917) {
918 tokio::spawn(async move {
919 match server.accept_request(session).await {
920 Ok(request) => {
921 let request = Request {
922 transport,
923 url: None,
924 identity: None,
925 kind: RequestKind::Qmux(Box::new(request)),
926 };
927 let _ = tx.send(request).await;
928 }
929 Err(err) => tracing::debug!(%err, "stream SETUP handshake failed"),
930 }
931 });
932}
933
934pub(crate) enum RequestKind {
941 #[cfg(feature = "noq")]
942 Noq(Box<moq_net::Request<web_transport_noq::Session>>),
943 #[cfg(feature = "quinn")]
944 Quinn(Box<moq_net::Request<web_transport_quinn::Session>>),
945 #[cfg(feature = "quiche")]
946 Quiche(Box<moq_net::Request<web_transport_quiche::Connection>>),
947 #[cfg(feature = "iroh")]
948 Iroh(Box<moq_net::Request<web_transport_iroh::Session>>),
949 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
950 Qmux(Box<moq_net::Request<qmux::Session>>),
951}
952
953#[non_exhaustive]
955#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
956pub enum Transport {
957 Quic,
959 Iroh,
961 WebSocket,
963 Tcp,
965 Unix,
967}
968
969impl Transport {
970 pub const fn as_str(self) -> &'static str {
972 match self {
973 Self::Quic => "quic",
974 Self::Iroh => "iroh",
975 Self::WebSocket => "websocket",
976 Self::Tcp => "tcp",
977 Self::Unix => "unix",
978 }
979 }
980}
981
982impl std::fmt::Display for Transport {
983 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
984 f.write_str(self.as_str())
985 }
986}
987
988pub struct Request {
997 transport: Transport,
998 url: Option<Url>,
1001 identity: Option<crate::tls::PeerIdentity>,
1004 kind: RequestKind,
1005}
1006
1007macro_rules! request_ref {
1009 ($self:expr, $r:ident => $body:expr) => {
1010 match &$self.kind {
1011 #[cfg(feature = "noq")]
1012 RequestKind::Noq($r) => $body,
1013 #[cfg(feature = "quinn")]
1014 RequestKind::Quinn($r) => $body,
1015 #[cfg(feature = "quiche")]
1016 RequestKind::Quiche($r) => $body,
1017 #[cfg(feature = "iroh")]
1018 RequestKind::Iroh($r) => $body,
1019 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
1020 RequestKind::Qmux($r) => $body,
1021 }
1022 };
1023}
1024
1025macro_rules! request_into {
1027 ($kind:expr, $r:ident => $body:expr) => {
1028 match $kind {
1029 #[cfg(feature = "noq")]
1030 RequestKind::Noq($r) => $body,
1031 #[cfg(feature = "quinn")]
1032 RequestKind::Quinn($r) => $body,
1033 #[cfg(feature = "quiche")]
1034 RequestKind::Quiche($r) => $body,
1035 #[cfg(feature = "iroh")]
1036 RequestKind::Iroh($r) => $body,
1037 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
1038 RequestKind::Qmux($r) => $body,
1039 }
1040 };
1041}
1042
1043macro_rules! request_map {
1045 ($kind:expr, $r:ident => $body:expr) => {
1046 match $kind {
1047 #[cfg(feature = "noq")]
1048 RequestKind::Noq($r) => RequestKind::Noq(Box::new($body)),
1049 #[cfg(feature = "quinn")]
1050 RequestKind::Quinn($r) => RequestKind::Quinn(Box::new($body)),
1051 #[cfg(feature = "quiche")]
1052 RequestKind::Quiche($r) => RequestKind::Quiche(Box::new($body)),
1053 #[cfg(feature = "iroh")]
1054 RequestKind::Iroh($r) => RequestKind::Iroh(Box::new($body)),
1055 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
1056 RequestKind::Qmux($r) => RequestKind::Qmux(Box::new($body)),
1057 }
1058 };
1059}
1060
1061impl Request {
1062 pub async fn close(self, code: u16) -> crate::Result<()> {
1066 let err = match code {
1067 401 | 403 => moq_net::Error::Unauthorized,
1068 other => moq_net::Error::App(other),
1069 };
1070 request_into!(self.kind, request => request.close(err));
1071 Ok(())
1072 }
1073
1074 pub fn with_publisher(self, publish: impl moq_net::Consume<moq_net::origin::Consumer>) -> Self {
1076 let Request {
1077 transport,
1078 url,
1079 identity,
1080 kind,
1081 } = self;
1082 let kind = request_map!(kind, request => request.with_publisher(publish));
1083 Request {
1084 transport,
1085 url,
1086 identity,
1087 kind,
1088 }
1089 }
1090
1091 pub fn with_subscriber(self, subscribe: moq_net::origin::Producer) -> Self {
1093 let Request {
1094 transport,
1095 url,
1096 identity,
1097 kind,
1098 } = self;
1099 let kind = request_map!(kind, request => request.with_subscriber(subscribe));
1100 Request {
1101 transport,
1102 url,
1103 identity,
1104 kind,
1105 }
1106 }
1107
1108 pub fn with_peer_origin(self, origin: moq_net::Origin) -> Self {
1112 let Request {
1113 transport,
1114 url,
1115 identity,
1116 kind,
1117 } = self;
1118 let kind = request_map!(kind, request => request.with_peer_origin(origin));
1119 Request {
1120 transport,
1121 url,
1122 identity,
1123 kind,
1124 }
1125 }
1126
1127 pub fn with_stats(self, stats: moq_net::stats::Session) -> Self {
1129 let Request {
1130 transport,
1131 url,
1132 identity,
1133 kind,
1134 } = self;
1135 let kind = request_map!(kind, request => request.with_stats(stats));
1136 Request {
1137 transport,
1138 url,
1139 identity,
1140 kind,
1141 }
1142 }
1143
1144 pub async fn ok(self) -> crate::Result<Session> {
1146 let pair = request_into!(self.kind, request => request.ok().await?);
1147 Ok(crate::spawn_session(pair))
1148 }
1149
1150 pub fn transport(&self) -> Transport {
1152 self.transport
1153 }
1154
1155 pub fn url(&self) -> Option<&Url> {
1160 self.url.as_ref()
1161 }
1162
1163 pub fn path(&self) -> &str {
1170 let setup = request_ref!(self, r => r.path());
1174 let path = if setup.is_empty() {
1175 self.url.as_ref().map(Url::path).unwrap_or("")
1176 } else {
1177 setup.split_once('?').map_or(setup, |(path, _)| path)
1178 };
1179 if path == "/" { "" } else { path }
1180 }
1181
1182 pub fn query(&self) -> Option<&str> {
1186 let setup = request_ref!(self, r => r.path());
1187 if setup.is_empty() {
1188 self.url.as_ref().and_then(Url::query)
1189 } else {
1190 setup.split_once('?').map(|(_, query)| query)
1191 }
1192 }
1193
1194 pub fn role(&self) -> Option<moq_net::Role> {
1199 request_ref!(self, r => r.role())
1200 }
1201
1202 pub fn peer_origin(&self) -> Option<moq_net::Origin> {
1210 request_ref!(self, r => r.peer_origin())
1211 }
1212
1213 pub fn peer_identity(&self) -> Option<crate::tls::PeerIdentity> {
1221 self.identity.clone()
1222 }
1223
1224 #[doc(hidden)]
1225 #[deprecated(note = "use `peer_identity` instead")]
1226 pub fn has_peer_certificate(&self) -> bool {
1227 self.peer_identity().is_some()
1228 }
1229}
1230
1231#[cfg(test)]
1232mod tests {
1233 use super::*;
1234
1235 #[test]
1236 fn version_help_lists_every_parseable_name() {
1237 let help = <ServerConfig as clap::Args>::augment_args(clap::Command::new("test"))
1238 .render_long_help()
1239 .to_string();
1240 for name in moq_net::Version::names() {
1241 assert!(help.contains(name), "missing {name} from --server-version help");
1242 }
1243 }
1244
1245 #[cfg(feature = "tcp")]
1253 #[test]
1254 fn accept_health_covers_stream_listeners_before_they_bind() {
1255 let mut config = ServerConfig::default();
1256 config.tcp.bind = Some("127.0.0.1:0".parse().unwrap());
1257 let server = Server::new(config).expect("stream-only server");
1258
1259 let names: Vec<_> = server.accept_health().iter().map(|h| h.listener()).collect();
1260 assert_eq!(names, vec!["tcp"], "the tcp listener must report before it binds");
1261 }
1262
1263 #[cfg(all(feature = "tcp", feature = "uds", unix))]
1269 #[tokio::test]
1270 async fn a_failed_listen_binds_nothing_and_can_be_retried() {
1271 let dir = tempfile::TempDir::new().unwrap();
1274 let occupied = dir.path().join("not-a-socket");
1275 std::fs::write(&occupied, b"in the way").unwrap();
1276
1277 let mut config = ServerConfig::default();
1278 config.tcp.bind = Some("127.0.0.1:0".parse().unwrap());
1279 config.unix.bind = Some(occupied);
1280 let mut server = Server::new(config).expect("stream-only server");
1281
1282 assert!(server.listen().await.is_err(), "the unix bind must fail");
1283 assert!(server.listen().await.is_err(), "a retry must not report success");
1286 }
1287
1288 #[cfg(feature = "tcp")]
1291 #[tokio::test]
1292 async fn close_releases_stream_listener_socket() {
1293 let probe = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1294 let addr = probe.local_addr().unwrap();
1295 drop(probe);
1296
1297 let mut config = ServerConfig::default();
1298 config.tcp.bind = Some(addr);
1299 let mut server = Server::new(config).expect("stream-only server");
1300 server.listen().await.expect("listen");
1301 assert!(tokio::net::TcpListener::bind(addr).await.is_err(), "listener is bound");
1302
1303 server.close().await;
1304 server.listen().await.expect("closed listener stays terminal");
1305 let _rebound = tokio::net::TcpListener::bind(addr)
1306 .await
1307 .expect("close must release the listener socket");
1308 }
1309
1310 #[cfg(not(any(feature = "noq", feature = "quinn", feature = "quiche")))]
1312 #[test]
1313 fn quic_bind_without_a_quic_backend_is_rejected() {
1314 let config = ServerConfig {
1315 bind: Some("127.0.0.1:0".to_string()),
1316 ..Default::default()
1317 };
1318
1319 assert!(matches!(Server::new(config), Err(Error::NoBackend(_))));
1320 }
1321
1322 #[cfg(all(feature = "quinn", not(feature = "tcp")))]
1325 #[test]
1326 fn accept_health_is_empty_without_a_stream_listener() {
1327 let server = ServerConfig::default().init().expect("quic server");
1328 assert!(server.accept_health().is_empty());
1329 }
1330
1331 #[test]
1332 fn transport_names_are_stable() {
1333 assert_eq!(Transport::Quic.as_str(), "quic");
1334 assert_eq!(Transport::Iroh.as_str(), "iroh");
1335 assert_eq!(Transport::WebSocket.as_str(), "websocket");
1336 assert_eq!(Transport::Tcp.as_str(), "tcp");
1337 assert_eq!(Transport::Unix.as_str(), "unix");
1338 }
1339
1340 #[cfg(feature = "quinn")]
1343 #[tokio::test]
1344 async fn certificates_expose_generated_fingerprints() {
1345 let mut config = ServerConfig {
1346 bind: Some("[::]:0".to_string()),
1347 ..Default::default()
1348 };
1349 config.tls.generate = vec!["localhost".into()];
1350
1351 let certs = config.init().expect("server init").certificates();
1352 let fingerprints = certs.fingerprints();
1353 assert_eq!(fingerprints.len(), 1, "one generated certificate");
1354 assert_eq!(fingerprints[0].len(), 64);
1356 assert!(fingerprints[0].chars().all(|c| c.is_ascii_hexdigit()));
1357 }
1358
1359 #[cfg(all(feature = "uds", unix))]
1364 #[tokio::test]
1365 async fn unix_listener_serves_the_configured_publisher() {
1366 use rand::RngExt;
1367
1368 let path = PathBuf::from(format!("/tmp/moq-native-publish-{}.sock", std::process::id()));
1371 let _ = std::fs::remove_file(&path);
1372
1373 let origin = moq_net::Origin::random().produce();
1374 let mut broadcast = origin
1375 .create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
1376 .expect("create broadcast");
1377 let mut track = broadcast.create_track("video", None).expect("create track");
1378 let mut group = track.append_group().expect("append group");
1379 group
1380 .write_frame(moq_net::Timestamp::ZERO, b"hello".as_ref())
1381 .expect("write frame");
1382 group.finish().expect("finish group");
1383
1384 let mut config = ServerConfig::default();
1385 config.unix.bind = Some(path.clone());
1386 let server = config.init().expect("server init");
1387
1388 let serve = tokio::spawn(server.serve_publish(origin.consume()));
1390
1391 const MAX_DELAY: std::time::Duration = std::time::Duration::from_millis(100);
1395 let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
1396 let mut delay = std::time::Duration::from_millis(1);
1397 while let Err(err) = tokio::net::UnixStream::connect(&path).await {
1398 assert!(
1399 tokio::time::Instant::now() < deadline,
1400 "unix listener never bound: {err}"
1401 );
1402 tokio::time::sleep(delay.mul_f64(0.5 + rand::rng().random::<f64>() / 2.0)).await;
1403 delay = (delay * 2).min(MAX_DELAY);
1404 }
1405
1406 const TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
1407
1408 let url: Url = format!("unix://{}", path.display()).parse().expect("parse url");
1409 let subscriber = moq_net::Origin::random().produce();
1410 let mut announced = subscriber.consume().announced();
1411 let client = crate::ClientConfig::default()
1412 .init()
1413 .expect("client init")
1414 .with_subscriber(subscriber);
1415 let session = tokio::time::timeout(TIMEOUT, client.connect(url))
1416 .await
1417 .expect("connect timeout")
1418 .expect("connect");
1419
1420 let update = tokio::time::timeout(TIMEOUT, announced.next())
1423 .await
1424 .expect("announce timeout")
1425 .expect("origin closed");
1426 assert_eq!(update.path.as_str(), "test");
1427 let broadcast = update.broadcast.expect("expected an announce");
1428
1429 let mut track = broadcast
1430 .track("video")
1431 .expect("track name")
1432 .subscribe(None)
1433 .await
1434 .expect("subscribe");
1435 let mut group = tokio::time::timeout(TIMEOUT, track.recv_group())
1436 .await
1437 .expect("recv group timeout")
1438 .expect("recv group")
1439 .expect("track closed early");
1440 let frame = tokio::time::timeout(TIMEOUT, group.read_frame())
1441 .await
1442 .expect("read frame timeout")
1443 .expect("read frame")
1444 .expect("group closed early");
1445 assert_eq!(&frame.payload[..], b"hello");
1446
1447 drop(session);
1448 serve.abort();
1449 let _ = std::fs::remove_file(&path);
1450 }
1451
1452 #[cfg(all(feature = "uds", unix))]
1455 #[tokio::test]
1456 async fn certificates_are_empty_without_a_tls_backend() {
1457 let mut config = ServerConfig::default();
1458 config.unix.bind = Some(PathBuf::from("/tmp/moq-native-certificates-test.sock"));
1459
1460 let server = config.init().expect("server init");
1461 assert!(server.certificates().fingerprints().is_empty());
1462 }
1463
1464 #[test]
1465 fn test_tls_string_or_array() {
1466 let single = r#"
1468 cert = "cert.pem"
1469 key = "key.pem"
1470 "#;
1471 let config: crate::tls::Server = toml::from_str(single).unwrap();
1472 assert_eq!(config.cert, vec![PathBuf::from("cert.pem")]);
1473 assert_eq!(config.key, vec![PathBuf::from("key.pem")]);
1474
1475 let array = r#"
1477 cert = ["a.pem", "b.pem"]
1478 key = ["a.key", "b.key"]
1479 generate = ["localhost"]
1480 root = ["ca.pem"]
1481 "#;
1482 let config: crate::tls::Server = toml::from_str(array).unwrap();
1483 assert_eq!(config.cert, vec![PathBuf::from("a.pem"), PathBuf::from("b.pem")]);
1484 assert_eq!(config.key, vec![PathBuf::from("a.key"), PathBuf::from("b.key")]);
1485 assert_eq!(config.generate, vec!["localhost".to_string()]);
1486 assert_eq!(config.root, vec![PathBuf::from("ca.pem")]);
1487 }
1488
1489 #[test]
1490 fn bind_string_or_listen_alias() {
1491 let bind: ServerConfig = toml::from_str(r#"bind = "[::]:443""#).unwrap();
1493 assert_eq!(bind.bind.as_deref(), Some("[::]:443"));
1494
1495 let alias: ServerConfig = toml::from_str(r#"listen = "0.0.0.0:4443""#).unwrap();
1496 assert_eq!(alias.bind.as_deref(), Some("0.0.0.0:4443"));
1497 }
1498
1499 #[cfg(all(feature = "uds", unix))]
1500 #[test]
1501 fn stream_listener_config_parses() {
1502 let config: ServerConfig = toml::from_str(
1503 r#"
1504bind = "[::]:443"
1505
1506[unix]
1507bind = "/run/moq.sock"
1508
1509[unix.allow]
1510uid = [1001, 1002]
1511"#,
1512 )
1513 .unwrap();
1514 assert_eq!(config.bind.as_deref(), Some("[::]:443"));
1515 assert_eq!(config.unix.bind.as_deref(), Some(std::path::Path::new("/run/moq.sock")));
1516 assert_eq!(config.unix.allow.as_ref().expect("allow").uid, vec![1001, 1002]);
1517 assert!(config.has_stream_listener());
1518 }
1519
1520 #[cfg(all(feature = "uds", unix))]
1521 #[test]
1522 fn stream_only_config_has_no_quic() {
1523 let mut config = ServerConfig::default();
1525 config.unix.bind = Some(PathBuf::from("/run/moq.sock"));
1526 assert!(config.has_stream_listener());
1527 assert!(config.bind.is_none());
1528
1529 assert!(!ServerConfig::default().has_stream_listener());
1531 }
1532}