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 async fn serve(mut self) -> crate::Result<()> {
314 if let Ok(addr) = self.local_addr() {
315 tracing::info!(%addr, "listening");
316 }
317 while let Some(request) = self.accept().await {
318 tokio::spawn(async move {
319 if let Err(err) = serve_session(request).await {
320 tracing::warn!(%err, "session ended with error");
321 }
322 });
323 }
324 Ok(())
325 }
326
327 pub fn certificates(&self) -> crate::tls::Certificates {
336 #[cfg(feature = "noq")]
337 if let Some(noq) = self.noq.as_ref() {
338 return noq.certificates();
339 }
340 #[cfg(feature = "quinn")]
341 if let Some(quinn) = self.quinn.as_ref() {
342 return quinn.certificates();
343 }
344 #[cfg(feature = "quiche")]
345 if let Some(quiche) = self.quiche.as_ref() {
346 return quiche.certificates();
347 }
348 crate::tls::Certificates::empty()
350 }
351
352 #[cfg(not(any(
353 feature = "noq",
354 feature = "quinn",
355 feature = "quiche",
356 feature = "iroh",
357 feature = "websocket",
358 feature = "tcp",
359 all(feature = "uds", unix)
360 )))]
361 pub async fn accept(&mut self) -> Option<Request> {
365 unreachable!("no transport compiled; enable a QUIC backend, websocket, tcp, or uds feature");
366 }
367
368 pub fn accept_health(&self) -> Vec<crate::accept::Health> {
380 #[allow(unused_mut)]
381 let mut health = Vec::new();
382 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
383 health.extend(self.streams.health.iter().cloned());
384 #[cfg(feature = "websocket")]
385 health.extend(self.websocket.as_ref().map(|ws| ws.accept_health()));
386 health
387 }
388
389 pub async fn listen(&mut self) -> crate::Result<()> {
406 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
407 self.streams.ensure_started(self.moq.clone()).await?;
408 Ok(())
409 }
410
411 #[cfg(any(
422 feature = "noq",
423 feature = "quinn",
424 feature = "quiche",
425 feature = "iroh",
426 feature = "websocket",
427 feature = "tcp",
428 all(feature = "uds", unix)
429 ))]
430 pub async fn accept(&mut self) -> Option<Request> {
431 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
435 if let Err(err) = self.streams.ensure_started(self.moq.clone()).await {
436 tracing::error!(%err, "failed to bind stream listener");
437 return None;
438 }
439
440 loop {
441 #[cfg(feature = "noq")]
443 let noq_accept = async {
444 #[cfg(feature = "noq")]
445 if let Some(noq) = self.noq.as_mut() {
446 return noq.accept().await;
447 }
448 None
449 };
450 #[cfg(not(feature = "noq"))]
451 let noq_accept = async { None::<()> };
452
453 #[cfg(feature = "iroh")]
454 let iroh_accept = async {
455 #[cfg(feature = "iroh")]
456 if let Some(endpoint) = self.iroh.as_mut() {
457 return endpoint.accept().await;
458 }
459 None
460 };
461 #[cfg(not(feature = "iroh"))]
462 let iroh_accept = async { None::<()> };
463
464 #[cfg(feature = "quinn")]
465 let quinn_accept = async {
466 #[cfg(feature = "quinn")]
467 if let Some(quinn) = self.quinn.as_mut() {
468 return quinn.accept().await;
469 }
470 None
471 };
472 #[cfg(not(feature = "quinn"))]
473 let quinn_accept = async { None::<()> };
474
475 #[cfg(feature = "quiche")]
476 let quiche_accept = async {
477 #[cfg(feature = "quiche")]
478 if let Some(quiche) = self.quiche.as_mut() {
479 return quiche.accept().await;
480 }
481 None
482 };
483 #[cfg(not(feature = "quiche"))]
484 let quiche_accept = async { None::<()> };
485
486 #[cfg(feature = "websocket")]
487 let ws_ref = self.websocket.as_ref();
488 #[cfg(feature = "websocket")]
489 let ws_accept = async {
490 match ws_ref {
491 Some(ws) => ws.accept_with_url().await,
492 None => std::future::pending().await,
493 }
494 };
495 #[cfg(not(feature = "websocket"))]
496 let ws_accept = std::future::pending::<Option<crate::Result<()>>>();
497
498 #[allow(unused_variables)]
499 let server = self.moq.clone();
500 #[allow(unused_variables)]
501 let versions = self.versions.clone();
502
503 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
505 let stream_accept = self.streams.recv();
506 #[cfg(not(any(feature = "tcp", all(feature = "uds", unix))))]
507 let stream_accept = std::future::pending::<Option<Request>>();
508
509 tokio::select! {
510 Some(request) = stream_accept => {
511 return Some(request);
512 }
513 Some(_conn) = noq_accept => {
514 #[cfg(feature = "noq")]
515 {
516 let alpns = versions.alpns();
517 self.accept.push(async move {
518 let (session, url, identity) = super::noq::accept(_conn, alpns).await?;
522 let request = server.accept_request(session).await?;
523 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Noq(Box::new(request)) })
524 }.boxed());
525 }
526 }
527 Some(_conn) = quinn_accept => {
528 #[cfg(feature = "quinn")]
529 {
530 let alpns = versions.alpns();
531 self.accept.push(async move {
532 let (session, url, identity) = super::quinn::accept(_conn, alpns).await?;
533 let request = server.accept_request(session).await?;
534 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Quinn(Box::new(request)) })
535 }.boxed());
536 }
537 }
538 Some(_conn) = quiche_accept => {
539 #[cfg(feature = "quiche")]
540 {
541 let alpns = versions.alpns();
542 self.accept.push(async move {
543 let (session, url, identity) = super::quiche::accept(_conn, alpns).await?;
544 let request = server.accept_request(session).await?;
545 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Quiche(Box::new(request)) })
546 }.boxed());
547 }
548 }
549 Some(_conn) = iroh_accept => {
550 #[cfg(feature = "iroh")]
551 self.accept.push(async move {
552 let (session, url, identity) = super::iroh::accept(_conn).await?;
553 let request = server.accept_request(session).await?;
554 Ok(Request { transport: Transport::Iroh, url, identity, kind: RequestKind::Iroh(Box::new(request)) })
555 }.boxed());
556 }
557 Some(_res) = ws_accept => {
558 #[cfg(feature = "websocket")]
559 match _res {
560 Ok((session, url)) => {
561 self.accept.push(async move {
564 let request = server.accept_request(session).await?;
565 Ok(Request { transport: Transport::WebSocket, url: Some(url), identity: None, kind: RequestKind::Qmux(Box::new(request)) })
566 }.boxed());
567 }
568 Err(err) => tracing::debug!(%err, "WebSocket upgrade failed"),
572 }
573 }
574 Some(res) = self.accept.next() => {
575 match res {
576 Ok(session) => return Some(session),
577 Err(err) => tracing::debug!(%err, "failed to accept session"),
578 }
579 }
580 _ = tokio::signal::ctrl_c() => {
581 self.close().await;
582 return None;
583 }
584 }
585 }
586 }
587
588 #[cfg(feature = "iroh")]
590 pub fn iroh_endpoint(&self) -> Option<&iroh::Endpoint> {
591 self.iroh.as_ref()
592 }
593
594 pub fn local_addr(&self) -> crate::Result<net::SocketAddr> {
600 #[cfg(feature = "noq")]
601 if let Some(noq) = self.noq.as_ref() {
602 return Ok(noq.local_addr()?);
603 }
604 #[cfg(feature = "quinn")]
605 if let Some(quinn) = self.quinn.as_ref() {
606 return Ok(quinn.local_addr()?);
607 }
608 #[cfg(feature = "quiche")]
609 if let Some(quiche) = self.quiche.as_ref() {
610 return Ok(quiche.local_addr()?);
611 }
612 Err(Error::NoBackend("no QUIC listener configured"))
614 }
615
616 #[cfg(feature = "websocket")]
619 pub fn websocket_local_addr(&self) -> Option<net::SocketAddr> {
620 self.websocket.as_ref().and_then(|ws| ws.local_addr().ok())
621 }
622
623 pub async fn close(&mut self) {
628 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
629 self.streams.close().await;
630 #[cfg(feature = "noq")]
631 if let Some(noq) = self.noq.as_mut() {
632 noq.close();
633 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
634 }
635 #[cfg(feature = "quinn")]
636 if let Some(quinn) = self.quinn.as_mut() {
637 quinn.close();
638 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
639 }
640 #[cfg(feature = "quiche")]
641 if let Some(quiche) = self.quiche.as_mut() {
642 quiche.close();
643 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
644 }
645 #[cfg(feature = "iroh")]
646 if let Some(iroh) = self.iroh.take() {
647 iroh.close().await;
648 }
649 #[cfg(feature = "websocket")]
650 {
651 let _ = self.websocket.take();
652 }
653 }
654}
655
656async fn serve_session(request: Request) -> crate::Result<()> {
658 let session = request.ok().await?;
659 Err(session.closed().await.into())
660}
661
662#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
668fn stream_versions(base: &moq_net::Versions) -> moq_net::Versions {
669 let mut versions: Vec<moq_net::Version> = base.iter().copied().collect();
670 if let Ok(lite05) = "moq-lite-05".parse::<moq_net::Version>()
671 && !versions.contains(&lite05)
672 {
673 versions.push(lite05);
674 }
675 moq_net::Versions::from(versions)
676}
677
678#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
680#[derive(Clone)]
681enum StreamBind {
682 #[cfg(feature = "tcp")]
683 Tcp(net::SocketAddr),
684 #[cfg(all(feature = "uds", unix))]
685 Unix(PathBuf),
686}
687
688#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
689impl StreamBind {
690 fn name(&self) -> &'static str {
692 match self {
693 #[cfg(feature = "tcp")]
694 Self::Tcp(_) => "tcp",
695 #[cfg(all(feature = "uds", unix))]
696 Self::Unix(_) => "unix",
697 }
698 }
699}
700
701#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
708struct StreamListeners {
709 binds: Vec<StreamBind>,
710 health: Vec<crate::accept::Health>,
714 versions: moq_net::Versions,
715 #[cfg(all(feature = "uds", unix))]
716 unix_allow: Option<crate::unix::Allow>,
717 rx: Option<tokio::sync::mpsc::Receiver<Request>>,
718 tasks: Vec<tokio::task::JoinHandle<()>>,
719}
720
721#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
722impl StreamListeners {
723 fn new(
724 binds: Vec<StreamBind>,
725 versions: moq_net::Versions,
726 #[cfg(all(feature = "uds", unix))] unix_allow: Option<crate::unix::Allow>,
727 ) -> Self {
728 let health = binds
729 .iter()
730 .map(|bind| crate::accept::Health::new(bind.name()))
731 .collect();
732 Self {
733 binds,
734 health,
735 versions,
736 #[cfg(all(feature = "uds", unix))]
737 unix_allow,
738 rx: None,
739 tasks: Vec::new(),
740 }
741 }
742
743 async fn ensure_started(&mut self, server: moq_net::Server) -> crate::Result<()> {
748 if self.rx.is_some() || self.binds.is_empty() {
749 return Ok(());
750 }
751
752 let server = server.with_versions(self.versions.clone());
755
756 let (tx, rx) = tokio::sync::mpsc::channel(16);
757 if let Err(err) = self.start(&server, &tx).await {
758 for task in self.tasks.drain(..) {
764 task.abort();
765 }
766 return Err(err);
767 }
768
769 self.rx = Some(rx);
770 Ok(())
771 }
772
773 async fn start(&mut self, server: &moq_net::Server, tx: &tokio::sync::mpsc::Sender<Request>) -> crate::Result<()> {
775 let binds = self.binds.clone();
778 let health = self.health.clone();
779 for (bind, health) in binds.into_iter().zip(health) {
780 let alpns = self.versions.alpns();
781 match bind {
782 #[cfg(feature = "tcp")]
783 StreamBind::Tcp(addr) => {
784 if !addr.ip().is_loopback() {
785 tracing::warn!(%addr, "tcp listener bound to a non-loopback address; qmux is UNENCRYPTED, ensure the network is trusted");
786 }
787 let listener = crate::tcp::Listener::bind(addr)
788 .await?
789 .with_protocols(alpns)
790 .with_accept_health(health);
791 tracing::info!(%addr, "listening (tcp)");
792 self.tasks.push(spawn_tcp_loop(listener, server.clone(), tx.clone()));
793 }
794 #[cfg(all(feature = "uds", unix))]
795 StreamBind::Unix(path) => {
796 let listener = crate::unix::Listener::bind(&path)
797 .await?
798 .with_protocols(alpns)
799 .with_accept_health(health);
800 listener.set_mode(0o666)?;
803 tracing::info!(path = %path.display(), allow = ?self.unix_allow, "listening (unix)");
804 self.tasks.push(spawn_unix_loop(
805 listener,
806 server.clone(),
807 self.unix_allow.clone(),
808 tx.clone(),
809 ));
810 }
811 }
812 }
813
814 Ok(())
815 }
816
817 async fn recv(&mut self) -> Option<Request> {
819 match self.rx.as_mut() {
820 Some(rx) => rx.recv().await,
821 None => std::future::pending().await,
822 }
823 }
824
825 async fn close(&mut self) {
827 self.binds.clear();
828 self.rx = None;
829 let tasks = std::mem::take(&mut self.tasks);
830 for task in tasks {
831 task.abort();
832 let _ = task.await;
833 }
834 }
835}
836
837#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
838impl Drop for StreamListeners {
839 fn drop(&mut self) {
840 for task in &self.tasks {
842 task.abort();
843 }
844 }
845}
846
847#[cfg(feature = "tcp")]
848fn spawn_tcp_loop(
849 listener: crate::tcp::Listener,
850 server: moq_net::Server,
851 tx: tokio::sync::mpsc::Sender<Request>,
852) -> tokio::task::JoinHandle<()> {
853 tokio::spawn(async move {
854 loop {
855 match listener.accept().await {
856 Some(Ok(session)) => spawn_stream_request(session, Transport::Tcp, server.clone(), tx.clone()),
857 Some(Err(err)) => tracing::warn!(%err, "tcp qmux handshake failed"),
860 None => break,
861 }
862 }
863 })
864}
865
866#[cfg(all(feature = "uds", unix))]
867fn spawn_unix_loop(
868 listener: crate::unix::Listener,
869 server: moq_net::Server,
870 allow: Option<crate::unix::Allow>,
871 tx: tokio::sync::mpsc::Sender<Request>,
872) -> tokio::task::JoinHandle<()> {
873 tokio::spawn(async move {
874 loop {
875 match listener.accept().await {
876 Some(Ok((session, cred))) => {
877 if let Some(allow) = &allow
879 && !allow.permits(&cred)
880 {
881 tracing::warn!(uid = cred.uid, gid = cred.gid, pid = ?cred.pid, "unix connection rejected by allow list");
882 continue;
883 }
884 spawn_stream_request(session, Transport::Unix, server.clone(), tx.clone());
885 }
886 Some(Err(err)) => tracing::warn!(%err, "unix qmux handshake failed"),
888 None => break,
889 }
890 }
891 })
892}
893
894#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
897fn spawn_stream_request(
898 session: qmux::Session,
899 transport: Transport,
900 server: moq_net::Server,
901 tx: tokio::sync::mpsc::Sender<Request>,
902) {
903 tokio::spawn(async move {
904 match server.accept_request(session).await {
905 Ok(request) => {
906 let request = Request {
907 transport,
908 url: None,
909 identity: None,
910 kind: RequestKind::Qmux(Box::new(request)),
911 };
912 let _ = tx.send(request).await;
913 }
914 Err(err) => tracing::debug!(%err, "stream SETUP handshake failed"),
915 }
916 });
917}
918
919pub(crate) enum RequestKind {
926 #[cfg(feature = "noq")]
927 Noq(Box<moq_net::Request<web_transport_noq::Session>>),
928 #[cfg(feature = "quinn")]
929 Quinn(Box<moq_net::Request<web_transport_quinn::Session>>),
930 #[cfg(feature = "quiche")]
931 Quiche(Box<moq_net::Request<web_transport_quiche::Connection>>),
932 #[cfg(feature = "iroh")]
933 Iroh(Box<moq_net::Request<web_transport_iroh::Session>>),
934 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
935 Qmux(Box<moq_net::Request<qmux::Session>>),
936}
937
938#[non_exhaustive]
940#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
941pub enum Transport {
942 Quic,
944 Iroh,
946 WebSocket,
948 Tcp,
950 Unix,
952}
953
954impl Transport {
955 pub const fn as_str(self) -> &'static str {
957 match self {
958 Self::Quic => "quic",
959 Self::Iroh => "iroh",
960 Self::WebSocket => "websocket",
961 Self::Tcp => "tcp",
962 Self::Unix => "unix",
963 }
964 }
965}
966
967impl std::fmt::Display for Transport {
968 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
969 f.write_str(self.as_str())
970 }
971}
972
973pub struct Request {
982 transport: Transport,
983 url: Option<Url>,
986 identity: Option<crate::tls::PeerIdentity>,
989 kind: RequestKind,
990}
991
992macro_rules! request_ref {
994 ($self:expr, $r:ident => $body:expr) => {
995 match &$self.kind {
996 #[cfg(feature = "noq")]
997 RequestKind::Noq($r) => $body,
998 #[cfg(feature = "quinn")]
999 RequestKind::Quinn($r) => $body,
1000 #[cfg(feature = "quiche")]
1001 RequestKind::Quiche($r) => $body,
1002 #[cfg(feature = "iroh")]
1003 RequestKind::Iroh($r) => $body,
1004 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
1005 RequestKind::Qmux($r) => $body,
1006 }
1007 };
1008}
1009
1010macro_rules! request_into {
1012 ($kind:expr, $r:ident => $body:expr) => {
1013 match $kind {
1014 #[cfg(feature = "noq")]
1015 RequestKind::Noq($r) => $body,
1016 #[cfg(feature = "quinn")]
1017 RequestKind::Quinn($r) => $body,
1018 #[cfg(feature = "quiche")]
1019 RequestKind::Quiche($r) => $body,
1020 #[cfg(feature = "iroh")]
1021 RequestKind::Iroh($r) => $body,
1022 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
1023 RequestKind::Qmux($r) => $body,
1024 }
1025 };
1026}
1027
1028macro_rules! request_map {
1030 ($kind:expr, $r:ident => $body:expr) => {
1031 match $kind {
1032 #[cfg(feature = "noq")]
1033 RequestKind::Noq($r) => RequestKind::Noq(Box::new($body)),
1034 #[cfg(feature = "quinn")]
1035 RequestKind::Quinn($r) => RequestKind::Quinn(Box::new($body)),
1036 #[cfg(feature = "quiche")]
1037 RequestKind::Quiche($r) => RequestKind::Quiche(Box::new($body)),
1038 #[cfg(feature = "iroh")]
1039 RequestKind::Iroh($r) => RequestKind::Iroh(Box::new($body)),
1040 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
1041 RequestKind::Qmux($r) => RequestKind::Qmux(Box::new($body)),
1042 }
1043 };
1044}
1045
1046impl Request {
1047 pub async fn close(self, code: u16) -> crate::Result<()> {
1051 let err = match code {
1052 401 | 403 => moq_net::Error::Unauthorized,
1053 other => moq_net::Error::App(other),
1054 };
1055 request_into!(self.kind, request => request.close(err));
1056 Ok(())
1057 }
1058
1059 pub fn with_publisher(self, publish: impl moq_net::Consume<moq_net::origin::Consumer>) -> Self {
1061 let Request {
1062 transport,
1063 url,
1064 identity,
1065 kind,
1066 } = self;
1067 let kind = request_map!(kind, request => request.with_publisher(publish));
1068 Request {
1069 transport,
1070 url,
1071 identity,
1072 kind,
1073 }
1074 }
1075
1076 pub fn with_subscriber(self, subscribe: moq_net::origin::Producer) -> Self {
1078 let Request {
1079 transport,
1080 url,
1081 identity,
1082 kind,
1083 } = self;
1084 let kind = request_map!(kind, request => request.with_subscriber(subscribe));
1085 Request {
1086 transport,
1087 url,
1088 identity,
1089 kind,
1090 }
1091 }
1092
1093 pub fn with_stats(self, stats: moq_net::stats::Session) -> Self {
1095 let Request {
1096 transport,
1097 url,
1098 identity,
1099 kind,
1100 } = self;
1101 let kind = request_map!(kind, request => request.with_stats(stats));
1102 Request {
1103 transport,
1104 url,
1105 identity,
1106 kind,
1107 }
1108 }
1109
1110 pub async fn ok(self) -> crate::Result<Session> {
1112 let pair = request_into!(self.kind, request => request.ok().await?);
1113 Ok(crate::spawn_session(pair))
1114 }
1115
1116 pub fn transport(&self) -> Transport {
1118 self.transport
1119 }
1120
1121 pub fn url(&self) -> Option<&Url> {
1126 self.url.as_ref()
1127 }
1128
1129 pub fn path(&self) -> &str {
1136 let setup = request_ref!(self, r => r.path());
1140 let path = if setup.is_empty() {
1141 self.url.as_ref().map(Url::path).unwrap_or("")
1142 } else {
1143 setup.split_once('?').map_or(setup, |(path, _)| path)
1144 };
1145 if path == "/" { "" } else { path }
1146 }
1147
1148 pub fn query(&self) -> Option<&str> {
1152 let setup = request_ref!(self, r => r.path());
1153 if setup.is_empty() {
1154 self.url.as_ref().and_then(Url::query)
1155 } else {
1156 setup.split_once('?').map(|(_, query)| query)
1157 }
1158 }
1159
1160 pub fn role(&self) -> Option<moq_net::Role> {
1165 request_ref!(self, r => r.role())
1166 }
1167
1168 pub fn peer_origin(&self) -> Option<moq_net::Origin> {
1176 request_ref!(self, r => r.peer_origin())
1177 }
1178
1179 pub fn peer_identity(&self) -> Option<crate::tls::PeerIdentity> {
1187 self.identity.clone()
1188 }
1189
1190 #[doc(hidden)]
1191 #[deprecated(note = "use `peer_identity` instead")]
1192 pub fn has_peer_certificate(&self) -> bool {
1193 self.peer_identity().is_some()
1194 }
1195}
1196
1197#[cfg(test)]
1198mod tests {
1199 use super::*;
1200
1201 #[test]
1202 fn version_help_lists_every_parseable_name() {
1203 let help = <ServerConfig as clap::Args>::augment_args(clap::Command::new("test"))
1204 .render_long_help()
1205 .to_string();
1206 for name in moq_net::Version::names() {
1207 assert!(help.contains(name), "missing {name} from --server-version help");
1208 }
1209 }
1210
1211 #[cfg(feature = "tcp")]
1219 #[test]
1220 fn accept_health_covers_stream_listeners_before_they_bind() {
1221 let mut config = ServerConfig::default();
1222 config.tcp.bind = Some("127.0.0.1:0".parse().unwrap());
1223 let server = Server::new(config).expect("stream-only server");
1224
1225 let names: Vec<_> = server.accept_health().iter().map(|h| h.listener()).collect();
1226 assert_eq!(names, vec!["tcp"], "the tcp listener must report before it binds");
1227 }
1228
1229 #[cfg(all(feature = "tcp", feature = "uds", unix))]
1235 #[tokio::test]
1236 async fn a_failed_listen_binds_nothing_and_can_be_retried() {
1237 let dir = tempfile::TempDir::new().unwrap();
1240 let occupied = dir.path().join("not-a-socket");
1241 std::fs::write(&occupied, b"in the way").unwrap();
1242
1243 let mut config = ServerConfig::default();
1244 config.tcp.bind = Some("127.0.0.1:0".parse().unwrap());
1245 config.unix.bind = Some(occupied);
1246 let mut server = Server::new(config).expect("stream-only server");
1247
1248 assert!(server.listen().await.is_err(), "the unix bind must fail");
1249 assert!(server.listen().await.is_err(), "a retry must not report success");
1252 }
1253
1254 #[cfg(feature = "tcp")]
1257 #[tokio::test]
1258 async fn close_releases_stream_listener_socket() {
1259 let probe = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1260 let addr = probe.local_addr().unwrap();
1261 drop(probe);
1262
1263 let mut config = ServerConfig::default();
1264 config.tcp.bind = Some(addr);
1265 let mut server = Server::new(config).expect("stream-only server");
1266 server.listen().await.expect("listen");
1267 assert!(tokio::net::TcpListener::bind(addr).await.is_err(), "listener is bound");
1268
1269 server.close().await;
1270 server.listen().await.expect("closed listener stays terminal");
1271 let _rebound = tokio::net::TcpListener::bind(addr)
1272 .await
1273 .expect("close must release the listener socket");
1274 }
1275
1276 #[cfg(not(any(feature = "noq", feature = "quinn", feature = "quiche")))]
1278 #[test]
1279 fn quic_bind_without_a_quic_backend_is_rejected() {
1280 let config = ServerConfig {
1281 bind: Some("127.0.0.1:0".to_string()),
1282 ..Default::default()
1283 };
1284
1285 assert!(matches!(Server::new(config), Err(Error::NoBackend(_))));
1286 }
1287
1288 #[cfg(all(feature = "quinn", not(feature = "tcp")))]
1291 #[test]
1292 fn accept_health_is_empty_without_a_stream_listener() {
1293 let server = ServerConfig::default().init().expect("quic server");
1294 assert!(server.accept_health().is_empty());
1295 }
1296
1297 #[test]
1298 fn transport_names_are_stable() {
1299 assert_eq!(Transport::Quic.as_str(), "quic");
1300 assert_eq!(Transport::Iroh.as_str(), "iroh");
1301 assert_eq!(Transport::WebSocket.as_str(), "websocket");
1302 assert_eq!(Transport::Tcp.as_str(), "tcp");
1303 assert_eq!(Transport::Unix.as_str(), "unix");
1304 }
1305
1306 #[cfg(feature = "quinn")]
1309 #[tokio::test]
1310 async fn certificates_expose_generated_fingerprints() {
1311 let mut config = ServerConfig {
1312 bind: Some("[::]:0".to_string()),
1313 ..Default::default()
1314 };
1315 config.tls.generate = vec!["localhost".into()];
1316
1317 let certs = config.init().expect("server init").certificates();
1318 let fingerprints = certs.fingerprints();
1319 assert_eq!(fingerprints.len(), 1, "one generated certificate");
1320 assert_eq!(fingerprints[0].len(), 64);
1322 assert!(fingerprints[0].chars().all(|c| c.is_ascii_hexdigit()));
1323 }
1324
1325 #[cfg(all(feature = "uds", unix))]
1330 #[tokio::test]
1331 async fn unix_listener_serves_the_configured_publisher() {
1332 use rand::RngExt;
1333
1334 let path = PathBuf::from(format!("/tmp/moq-native-publish-{}.sock", std::process::id()));
1337 let _ = std::fs::remove_file(&path);
1338
1339 let origin = moq_net::Origin::random().produce();
1340 let mut broadcast = origin
1341 .create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
1342 .expect("create broadcast");
1343 let mut track = broadcast.create_track("video", None).expect("create track");
1344 let mut group = track.append_group().expect("append group");
1345 group
1346 .write_frame(moq_net::Timestamp::ZERO, b"hello".as_ref())
1347 .expect("write frame");
1348 group.finish().expect("finish group");
1349
1350 let mut config = ServerConfig::default();
1351 config.unix.bind = Some(path.clone());
1352 let server = config.init().expect("server init");
1353
1354 let serve = tokio::spawn(server.serve_publish(origin.consume()));
1356
1357 const MAX_DELAY: std::time::Duration = std::time::Duration::from_millis(100);
1361 let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
1362 let mut delay = std::time::Duration::from_millis(1);
1363 while let Err(err) = tokio::net::UnixStream::connect(&path).await {
1364 assert!(
1365 tokio::time::Instant::now() < deadline,
1366 "unix listener never bound: {err}"
1367 );
1368 tokio::time::sleep(delay.mul_f64(0.5 + rand::rng().random::<f64>() / 2.0)).await;
1369 delay = (delay * 2).min(MAX_DELAY);
1370 }
1371
1372 const TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
1373
1374 let url: Url = format!("unix://{}", path.display()).parse().expect("parse url");
1375 let subscriber = moq_net::Origin::random().produce();
1376 let mut announced = subscriber.consume().announced();
1377 let client = crate::ClientConfig::default()
1378 .init()
1379 .expect("client init")
1380 .with_subscriber(subscriber);
1381 let session = tokio::time::timeout(TIMEOUT, client.connect(url))
1382 .await
1383 .expect("connect timeout")
1384 .expect("connect");
1385
1386 let update = tokio::time::timeout(TIMEOUT, announced.next())
1389 .await
1390 .expect("announce timeout")
1391 .expect("origin closed");
1392 assert_eq!(update.path.as_str(), "test");
1393 let broadcast = update.broadcast.expect("expected an announce");
1394
1395 let mut track = broadcast
1396 .track("video")
1397 .expect("track name")
1398 .subscribe(None)
1399 .await
1400 .expect("subscribe");
1401 let mut group = tokio::time::timeout(TIMEOUT, track.recv_group())
1402 .await
1403 .expect("recv group timeout")
1404 .expect("recv group")
1405 .expect("track closed early");
1406 let frame = tokio::time::timeout(TIMEOUT, group.read_frame())
1407 .await
1408 .expect("read frame timeout")
1409 .expect("read frame")
1410 .expect("group closed early");
1411 assert_eq!(&frame.payload[..], b"hello");
1412
1413 drop(session);
1414 serve.abort();
1415 let _ = std::fs::remove_file(&path);
1416 }
1417
1418 #[cfg(all(feature = "uds", unix))]
1421 #[tokio::test]
1422 async fn certificates_are_empty_without_a_tls_backend() {
1423 let mut config = ServerConfig::default();
1424 config.unix.bind = Some(PathBuf::from("/tmp/moq-native-certificates-test.sock"));
1425
1426 let server = config.init().expect("server init");
1427 assert!(server.certificates().fingerprints().is_empty());
1428 }
1429
1430 #[test]
1431 fn test_tls_string_or_array() {
1432 let single = r#"
1434 cert = "cert.pem"
1435 key = "key.pem"
1436 "#;
1437 let config: crate::tls::Server = toml::from_str(single).unwrap();
1438 assert_eq!(config.cert, vec![PathBuf::from("cert.pem")]);
1439 assert_eq!(config.key, vec![PathBuf::from("key.pem")]);
1440
1441 let array = r#"
1443 cert = ["a.pem", "b.pem"]
1444 key = ["a.key", "b.key"]
1445 generate = ["localhost"]
1446 root = ["ca.pem"]
1447 "#;
1448 let config: crate::tls::Server = toml::from_str(array).unwrap();
1449 assert_eq!(config.cert, vec![PathBuf::from("a.pem"), PathBuf::from("b.pem")]);
1450 assert_eq!(config.key, vec![PathBuf::from("a.key"), PathBuf::from("b.key")]);
1451 assert_eq!(config.generate, vec!["localhost".to_string()]);
1452 assert_eq!(config.root, vec![PathBuf::from("ca.pem")]);
1453 }
1454
1455 #[test]
1456 fn bind_string_or_listen_alias() {
1457 let bind: ServerConfig = toml::from_str(r#"bind = "[::]:443""#).unwrap();
1459 assert_eq!(bind.bind.as_deref(), Some("[::]:443"));
1460
1461 let alias: ServerConfig = toml::from_str(r#"listen = "0.0.0.0:4443""#).unwrap();
1462 assert_eq!(alias.bind.as_deref(), Some("0.0.0.0:4443"));
1463 }
1464
1465 #[cfg(all(feature = "uds", unix))]
1466 #[test]
1467 fn stream_listener_config_parses() {
1468 let config: ServerConfig = toml::from_str(
1469 r#"
1470bind = "[::]:443"
1471
1472[unix]
1473bind = "/run/moq.sock"
1474
1475[unix.allow]
1476uid = [1001, 1002]
1477"#,
1478 )
1479 .unwrap();
1480 assert_eq!(config.bind.as_deref(), Some("[::]:443"));
1481 assert_eq!(config.unix.bind.as_deref(), Some(std::path::Path::new("/run/moq.sock")));
1482 assert_eq!(config.unix.allow.as_ref().expect("allow").uid, vec![1001, 1002]);
1483 assert!(config.has_stream_listener());
1484 }
1485
1486 #[cfg(all(feature = "uds", unix))]
1487 #[test]
1488 fn stream_only_config_has_no_quic() {
1489 let mut config = ServerConfig::default();
1491 config.unix.bind = Some(PathBuf::from("/run/moq.sock"));
1492 assert!(config.has_stream_listener());
1493 assert!(config.bind.is_none());
1494
1495 assert!(!ServerConfig::default().has_stream_listener());
1497 }
1498}