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
11use futures::FutureExt;
12use futures::future::BoxFuture;
13use futures::stream::FuturesUnordered;
14use futures::stream::StreamExt;
15
16#[derive(clap::Args, Clone, Debug, Default, serde::Serialize, serde::Deserialize)]
18#[serde(deny_unknown_fields, default)]
19#[non_exhaustive]
20pub struct ServerConfig {
21 #[serde(alias = "listen")]
29 #[arg(id = "server-bind", long = "server-bind", alias = "listen", env = "MOQ_SERVER_BIND")]
30 pub bind: Option<String>,
31
32 #[cfg(feature = "tcp")]
35 #[command(flatten)]
36 #[serde(default)]
37 pub tcp: crate::tcp::Config,
38
39 #[cfg(all(feature = "uds", unix))]
42 #[command(flatten)]
43 #[serde(default)]
44 pub unix: crate::unix::Config,
45
46 #[arg(id = "server-backend", long = "server-backend", env = "MOQ_SERVER_BACKEND")]
49 pub backend: Option<QuicBackend>,
50
51 #[command(flatten)]
54 #[serde(default)]
55 pub quic: crate::quic::Server,
56
57 #[serde(default, skip_serializing_if = "Vec::is_empty")]
65 #[arg(id = "server-version", long = "server-version", env = "MOQ_SERVER_VERSION")]
66 pub version: Vec<moq_net::Version>,
67
68 #[command(flatten)]
71 #[serde(default)]
72 pub tls: crate::tls::Server,
73}
74
75impl ServerConfig {
76 pub fn init(self) -> crate::Result<Server> {
78 Server::new(self)
79 }
80
81 pub fn versions(&self) -> moq_net::Versions {
83 if self.version.is_empty() {
84 moq_net::Versions::all()
85 } else {
86 moq_net::Versions::from(self.version.clone())
87 }
88 }
89
90 #[allow(unused_mut)]
95 fn has_stream_listener(&self) -> bool {
96 let mut has = false;
97 #[cfg(feature = "tcp")]
98 {
99 has |= self.tcp.bind.is_some();
100 }
101 #[cfg(all(feature = "uds", unix))]
102 {
103 has |= self.unix.bind.is_some();
104 }
105 has
106 }
107}
108
109pub(crate) const DEFAULT_BIND: &str = "[::]:443";
111
112pub struct Server {
118 moq: moq_net::Server,
119 versions: moq_net::Versions,
120 accept: FuturesUnordered<BoxFuture<'static, crate::Result<Request>>>,
121 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
122 streams: StreamListeners,
123 #[cfg(feature = "iroh")]
124 iroh: Option<iroh::Endpoint>,
125 #[cfg(feature = "noq")]
126 noq: Option<crate::noq::NoqServer>,
127 #[cfg(feature = "quinn")]
128 quinn: Option<crate::quinn::QuinnServer>,
129 #[cfg(feature = "quiche")]
130 quiche: Option<crate::quiche::QuicheServer>,
131 #[cfg(feature = "websocket")]
132 websocket: Option<crate::websocket::Listener>,
133}
134
135impl Server {
136 pub fn new(config: ServerConfig) -> crate::Result<Self> {
141 let backend = config.backend.clone().unwrap_or_else(crate::default_quic_backend);
142
143 let versions = config.versions();
144
145 config.quic.validate()?;
149
150 let build_quic = config.bind.is_some() || !config.has_stream_listener();
151
152 if build_quic && !config.tls.root.is_empty() {
153 let mtls_supported = match backend {
154 #[cfg(feature = "quinn")]
155 QuicBackend::Quinn => true,
156 #[cfg(feature = "noq")]
157 QuicBackend::Noq => true,
158 #[cfg(feature = "quiche")]
159 QuicBackend::Quiche => true,
160 #[allow(unreachable_patterns)]
161 _ => false,
162 };
163 if !mtls_supported {
164 return Err(Error::MtlsUnsupported);
165 }
166 }
167
168 #[cfg(feature = "noq")]
169 #[allow(unreachable_patterns)]
170 let noq = match backend {
171 QuicBackend::Noq if build_quic => Some(crate::noq::NoqServer::new(config.clone())?),
172 _ => None,
173 };
174
175 #[cfg(feature = "quinn")]
176 #[allow(unreachable_patterns)]
177 let quinn = match backend {
178 QuicBackend::Quinn if build_quic => Some(crate::quinn::QuinnServer::new(config.clone())?),
179 _ => None,
180 };
181
182 #[cfg(feature = "quiche")]
183 let quiche = match backend {
184 QuicBackend::Quiche if build_quic => Some(crate::quiche::QuicheServer::new(config.clone())?),
185 _ => None,
186 };
187
188 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
190 let mut stream_binds = Vec::new();
191 #[cfg(feature = "tcp")]
192 if let Some(addr) = config.tcp.bind {
193 stream_binds.push(StreamBind::Tcp(addr));
194 }
195 #[cfg(all(feature = "uds", unix))]
196 if let Some(path) = config.unix.bind.clone() {
197 stream_binds.push(StreamBind::Unix(path));
198 }
199 #[cfg(all(feature = "uds", unix))]
201 let unix_allow = config.unix.allow.clone().filter(|allow| !allow.is_empty());
202 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
203 let streams = StreamListeners::new(
204 stream_binds,
205 stream_versions(&versions),
206 #[cfg(all(feature = "uds", unix))]
207 unix_allow,
208 );
209
210 Ok(Server {
211 accept: Default::default(),
212 moq: moq_net::Server::new().with_versions(versions.clone()),
213 versions,
214 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
215 streams,
216 #[cfg(feature = "iroh")]
217 iroh: None,
218 #[cfg(feature = "noq")]
219 noq,
220 #[cfg(feature = "quinn")]
221 quinn,
222 #[cfg(feature = "quiche")]
223 quiche,
224 #[cfg(feature = "websocket")]
225 websocket: None,
226 })
227 }
228
229 #[cfg(feature = "websocket")]
235 pub fn with_websocket(mut self, websocket: crate::websocket::Listener) -> Self {
236 self.websocket = Some(websocket);
237 self
238 }
239
240 #[cfg(feature = "iroh")]
242 pub fn with_iroh(mut self, iroh: iroh::Endpoint) -> Self {
243 self.iroh = Some(iroh);
244 self
245 }
246
247 pub fn with_publisher(mut self, publish: impl moq_net::Consume<moq_net::origin::Consumer>) -> Self {
249 self.moq = self.moq.with_publisher(publish);
250 self
251 }
252
253 pub fn with_subscriber(mut self, subscribe: moq_net::origin::Producer) -> Self {
255 self.moq = self.moq.with_subscriber(subscribe);
256 self
257 }
258
259 pub fn with_stats(mut self, stats: moq_net::stats::Session) -> Self {
262 self.moq = self.moq.with_stats(stats);
263 self
264 }
265
266 pub async fn serve_publish(self, origin: moq_net::origin::Consumer) -> crate::Result<()> {
273 self.with_publisher(origin).serve().await
274 }
275
276 pub async fn serve_consume(self, origin: moq_net::origin::Producer) -> crate::Result<()> {
280 self.with_subscriber(origin).serve().await
281 }
282
283 async fn serve(mut self) -> crate::Result<()> {
286 if let Ok(addr) = self.local_addr() {
287 tracing::info!(%addr, "listening");
288 }
289 while let Some(request) = self.accept().await {
290 tokio::spawn(async move {
291 if let Err(err) = serve_session(request).await {
292 tracing::warn!(%err, "session ended with error");
293 }
294 });
295 }
296 Ok(())
297 }
298
299 pub fn certificates(&self) -> crate::tls::Certificates {
308 #[cfg(feature = "noq")]
309 if let Some(noq) = self.noq.as_ref() {
310 return noq.certificates();
311 }
312 #[cfg(feature = "quinn")]
313 if let Some(quinn) = self.quinn.as_ref() {
314 return quinn.certificates();
315 }
316 #[cfg(feature = "quiche")]
317 if let Some(quiche) = self.quiche.as_ref() {
318 return quiche.certificates();
319 }
320 crate::tls::Certificates::empty()
322 }
323
324 #[cfg(not(any(
325 feature = "noq",
326 feature = "quinn",
327 feature = "quiche",
328 feature = "iroh",
329 feature = "tcp",
330 all(feature = "uds", unix)
331 )))]
332 pub async fn accept(&mut self) -> Option<Request> {
336 unreachable!("no transport compiled; enable a QUIC backend, tcp, or uds feature");
337 }
338
339 pub fn accept_health(&self) -> Vec<crate::accept::Health> {
351 #[allow(unused_mut)]
352 let mut health = Vec::new();
353 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
354 health.extend(self.streams.health.iter().cloned());
355 #[cfg(feature = "websocket")]
356 health.extend(self.websocket.as_ref().map(|ws| ws.accept_health()));
357 health
358 }
359
360 pub async fn listen(&mut self) -> crate::Result<()> {
377 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
378 self.streams.ensure_started(self.moq.clone()).await?;
379 Ok(())
380 }
381
382 #[cfg(any(
393 feature = "noq",
394 feature = "quinn",
395 feature = "quiche",
396 feature = "iroh",
397 feature = "tcp",
398 all(feature = "uds", unix)
399 ))]
400 pub async fn accept(&mut self) -> Option<Request> {
401 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
405 if let Err(err) = self.streams.ensure_started(self.moq.clone()).await {
406 tracing::error!(%err, "failed to bind stream listener");
407 return None;
408 }
409
410 loop {
411 #[cfg(feature = "noq")]
413 let noq_accept = async {
414 #[cfg(feature = "noq")]
415 if let Some(noq) = self.noq.as_mut() {
416 return noq.accept().await;
417 }
418 None
419 };
420 #[cfg(not(feature = "noq"))]
421 let noq_accept = async { None::<()> };
422
423 #[cfg(feature = "iroh")]
424 let iroh_accept = async {
425 #[cfg(feature = "iroh")]
426 if let Some(endpoint) = self.iroh.as_mut() {
427 return endpoint.accept().await;
428 }
429 None
430 };
431 #[cfg(not(feature = "iroh"))]
432 let iroh_accept = async { None::<()> };
433
434 #[cfg(feature = "quinn")]
435 let quinn_accept = async {
436 #[cfg(feature = "quinn")]
437 if let Some(quinn) = self.quinn.as_mut() {
438 return quinn.accept().await;
439 }
440 None
441 };
442 #[cfg(not(feature = "quinn"))]
443 let quinn_accept = async { None::<()> };
444
445 #[cfg(feature = "quiche")]
446 let quiche_accept = async {
447 #[cfg(feature = "quiche")]
448 if let Some(quiche) = self.quiche.as_mut() {
449 return quiche.accept().await;
450 }
451 None
452 };
453 #[cfg(not(feature = "quiche"))]
454 let quiche_accept = async { None::<()> };
455
456 #[cfg(feature = "websocket")]
457 let ws_ref = self.websocket.as_ref();
458 #[cfg(feature = "websocket")]
459 let ws_accept = async {
460 match ws_ref {
461 Some(ws) => ws.accept().await,
462 None => std::future::pending().await,
463 }
464 };
465 #[cfg(not(feature = "websocket"))]
466 let ws_accept = std::future::pending::<Option<crate::Result<()>>>();
467
468 #[allow(unused_variables)]
469 let server = self.moq.clone();
470 #[allow(unused_variables)]
471 let versions = self.versions.clone();
472
473 #[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
475 let stream_accept = self.streams.recv();
476 #[cfg(not(any(feature = "tcp", all(feature = "uds", unix))))]
477 let stream_accept = std::future::pending::<Option<Request>>();
478
479 tokio::select! {
480 Some(request) = stream_accept => {
481 return Some(request);
482 }
483 Some(_conn) = noq_accept => {
484 #[cfg(feature = "noq")]
485 {
486 let alpns = versions.alpns();
487 self.accept.push(async move {
488 let (session, url, identity) = super::noq::accept(_conn, alpns).await?;
492 let request = server.accept_request(session).await?;
493 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Noq(Box::new(request)) })
494 }.boxed());
495 }
496 }
497 Some(_conn) = quinn_accept => {
498 #[cfg(feature = "quinn")]
499 {
500 let alpns = versions.alpns();
501 self.accept.push(async move {
502 let (session, url, identity) = super::quinn::accept(_conn, alpns).await?;
503 let request = server.accept_request(session).await?;
504 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Quinn(Box::new(request)) })
505 }.boxed());
506 }
507 }
508 Some(_conn) = quiche_accept => {
509 #[cfg(feature = "quiche")]
510 {
511 let alpns = versions.alpns();
512 self.accept.push(async move {
513 let (session, url, identity) = super::quiche::accept(_conn, alpns).await?;
514 let request = server.accept_request(session).await?;
515 Ok(Request { transport: Transport::Quic, url, identity, kind: RequestKind::Quiche(Box::new(request)) })
516 }.boxed());
517 }
518 }
519 Some(_conn) = iroh_accept => {
520 #[cfg(feature = "iroh")]
521 self.accept.push(async move {
522 let (session, url, identity) = super::iroh::accept(_conn).await?;
523 let request = server.accept_request(session).await?;
524 Ok(Request { transport: Transport::Iroh, url, identity, kind: RequestKind::Iroh(Box::new(request)) })
525 }.boxed());
526 }
527 Some(_res) = ws_accept => {
528 #[cfg(feature = "websocket")]
529 match _res {
530 Ok(session) => {
531 self.accept.push(async move {
534 let request = server.accept_request(session).await?;
535 Ok(Request { transport: Transport::WebSocket, url: None, identity: None, kind: RequestKind::Qmux(Box::new(request)) })
536 }.boxed());
537 }
538 Err(err) => tracing::debug!(%err, "WebSocket upgrade failed"),
542 }
543 }
544 Some(res) = self.accept.next() => {
545 match res {
546 Ok(session) => return Some(session),
547 Err(err) => tracing::debug!(%err, "failed to accept session"),
548 }
549 }
550 _ = tokio::signal::ctrl_c() => {
551 self.close().await;
552 return None;
553 }
554 }
555 }
556 }
557
558 #[cfg(feature = "iroh")]
560 pub fn iroh_endpoint(&self) -> Option<&iroh::Endpoint> {
561 self.iroh.as_ref()
562 }
563
564 pub fn local_addr(&self) -> crate::Result<net::SocketAddr> {
570 #[cfg(feature = "noq")]
571 if let Some(noq) = self.noq.as_ref() {
572 return Ok(noq.local_addr()?);
573 }
574 #[cfg(feature = "quinn")]
575 if let Some(quinn) = self.quinn.as_ref() {
576 return Ok(quinn.local_addr()?);
577 }
578 #[cfg(feature = "quiche")]
579 if let Some(quiche) = self.quiche.as_ref() {
580 return Ok(quiche.local_addr()?);
581 }
582 Err(Error::NoBackend("no QUIC listener configured"))
584 }
585
586 #[cfg(feature = "websocket")]
589 pub fn websocket_local_addr(&self) -> Option<net::SocketAddr> {
590 self.websocket.as_ref().and_then(|ws| ws.local_addr().ok())
591 }
592
593 pub async fn close(&mut self) {
598 #[cfg(feature = "noq")]
599 if let Some(noq) = self.noq.as_mut() {
600 noq.close();
601 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
602 }
603 #[cfg(feature = "quinn")]
604 if let Some(quinn) = self.quinn.as_mut() {
605 quinn.close();
606 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
607 }
608 #[cfg(feature = "quiche")]
609 if let Some(quiche) = self.quiche.as_mut() {
610 quiche.close();
611 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
612 }
613 #[cfg(feature = "iroh")]
614 if let Some(iroh) = self.iroh.take() {
615 iroh.close().await;
616 }
617 #[cfg(feature = "websocket")]
618 {
619 let _ = self.websocket.take();
620 }
621 #[cfg(not(any(feature = "noq", feature = "quinn", feature = "quiche", feature = "iroh")))]
622 unreachable!("no QUIC backend compiled");
623 }
624}
625
626async fn serve_session(request: Request) -> crate::Result<()> {
628 let session = request.ok().await?;
629 Err(session.closed().await.into())
630}
631
632#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
638fn stream_versions(base: &moq_net::Versions) -> moq_net::Versions {
639 let mut versions: Vec<moq_net::Version> = base.iter().copied().collect();
640 if let Ok(lite05) = "moq-lite-05".parse::<moq_net::Version>()
641 && !versions.contains(&lite05)
642 {
643 versions.push(lite05);
644 }
645 moq_net::Versions::from(versions)
646}
647
648#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
650#[derive(Clone)]
651enum StreamBind {
652 #[cfg(feature = "tcp")]
653 Tcp(net::SocketAddr),
654 #[cfg(all(feature = "uds", unix))]
655 Unix(PathBuf),
656}
657
658#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
659impl StreamBind {
660 fn name(&self) -> &'static str {
662 match self {
663 #[cfg(feature = "tcp")]
664 Self::Tcp(_) => "tcp",
665 #[cfg(all(feature = "uds", unix))]
666 Self::Unix(_) => "unix",
667 }
668 }
669}
670
671#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
678struct StreamListeners {
679 binds: Vec<StreamBind>,
680 health: Vec<crate::accept::Health>,
684 versions: moq_net::Versions,
685 #[cfg(all(feature = "uds", unix))]
686 unix_allow: Option<crate::unix::Allow>,
687 rx: Option<tokio::sync::mpsc::Receiver<Request>>,
688 tasks: Vec<tokio::task::JoinHandle<()>>,
689}
690
691#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
692impl StreamListeners {
693 fn new(
694 binds: Vec<StreamBind>,
695 versions: moq_net::Versions,
696 #[cfg(all(feature = "uds", unix))] unix_allow: Option<crate::unix::Allow>,
697 ) -> Self {
698 let health = binds
699 .iter()
700 .map(|bind| crate::accept::Health::new(bind.name()))
701 .collect();
702 Self {
703 binds,
704 health,
705 versions,
706 #[cfg(all(feature = "uds", unix))]
707 unix_allow,
708 rx: None,
709 tasks: Vec::new(),
710 }
711 }
712
713 async fn ensure_started(&mut self, server: moq_net::Server) -> crate::Result<()> {
718 if self.rx.is_some() || self.binds.is_empty() {
719 return Ok(());
720 }
721
722 let server = server.with_versions(self.versions.clone());
725
726 let (tx, rx) = tokio::sync::mpsc::channel(16);
727 if let Err(err) = self.start(&server, &tx).await {
728 for task in self.tasks.drain(..) {
734 task.abort();
735 }
736 return Err(err);
737 }
738
739 self.rx = Some(rx);
740 Ok(())
741 }
742
743 async fn start(&mut self, server: &moq_net::Server, tx: &tokio::sync::mpsc::Sender<Request>) -> crate::Result<()> {
745 let binds = self.binds.clone();
748 let health = self.health.clone();
749 for (bind, health) in binds.into_iter().zip(health) {
750 let alpns = self.versions.alpns();
751 match bind {
752 #[cfg(feature = "tcp")]
753 StreamBind::Tcp(addr) => {
754 if !addr.ip().is_loopback() {
755 tracing::warn!(%addr, "tcp listener bound to a non-loopback address; qmux is UNENCRYPTED, ensure the network is trusted");
756 }
757 let listener = crate::tcp::Listener::bind(addr)
758 .await?
759 .with_protocols(alpns)
760 .with_accept_health(health);
761 tracing::info!(%addr, "listening (tcp)");
762 self.tasks.push(spawn_tcp_loop(listener, server.clone(), tx.clone()));
763 }
764 #[cfg(all(feature = "uds", unix))]
765 StreamBind::Unix(path) => {
766 let listener = crate::unix::Listener::bind(&path)
767 .await?
768 .with_protocols(alpns)
769 .with_accept_health(health);
770 listener.set_mode(0o666)?;
773 tracing::info!(path = %path.display(), allow = ?self.unix_allow, "listening (unix)");
774 self.tasks.push(spawn_unix_loop(
775 listener,
776 server.clone(),
777 self.unix_allow.clone(),
778 tx.clone(),
779 ));
780 }
781 }
782 }
783
784 Ok(())
785 }
786
787 async fn recv(&mut self) -> Option<Request> {
789 match self.rx.as_mut() {
790 Some(rx) => rx.recv().await,
791 None => std::future::pending().await,
792 }
793 }
794}
795
796#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
797impl Drop for StreamListeners {
798 fn drop(&mut self) {
799 for task in &self.tasks {
801 task.abort();
802 }
803 }
804}
805
806#[cfg(feature = "tcp")]
807fn spawn_tcp_loop(
808 listener: crate::tcp::Listener,
809 server: moq_net::Server,
810 tx: tokio::sync::mpsc::Sender<Request>,
811) -> tokio::task::JoinHandle<()> {
812 tokio::spawn(async move {
813 loop {
814 match listener.accept().await {
815 Some(Ok(session)) => spawn_stream_request(session, Transport::Tcp, server.clone(), tx.clone()),
816 Some(Err(err)) => tracing::warn!(%err, "tcp qmux handshake failed"),
819 None => break,
820 }
821 }
822 })
823}
824
825#[cfg(all(feature = "uds", unix))]
826fn spawn_unix_loop(
827 listener: crate::unix::Listener,
828 server: moq_net::Server,
829 allow: Option<crate::unix::Allow>,
830 tx: tokio::sync::mpsc::Sender<Request>,
831) -> tokio::task::JoinHandle<()> {
832 tokio::spawn(async move {
833 loop {
834 match listener.accept().await {
835 Some(Ok((session, cred))) => {
836 if let Some(allow) = &allow
838 && !allow.permits(&cred)
839 {
840 tracing::warn!(uid = cred.uid, gid = cred.gid, pid = ?cred.pid, "unix connection rejected by allow list");
841 continue;
842 }
843 spawn_stream_request(session, Transport::Unix, server.clone(), tx.clone());
844 }
845 Some(Err(err)) => tracing::warn!(%err, "unix qmux handshake failed"),
847 None => break,
848 }
849 }
850 })
851}
852
853#[cfg(any(feature = "tcp", all(feature = "uds", unix)))]
856fn spawn_stream_request(
857 session: qmux::Session,
858 transport: Transport,
859 server: moq_net::Server,
860 tx: tokio::sync::mpsc::Sender<Request>,
861) {
862 tokio::spawn(async move {
863 match server.accept_request(session).await {
864 Ok(request) => {
865 let request = Request {
866 transport,
867 url: None,
868 identity: None,
869 kind: RequestKind::Qmux(Box::new(request)),
870 };
871 let _ = tx.send(request).await;
872 }
873 Err(err) => tracing::debug!(%err, "stream SETUP handshake failed"),
874 }
875 });
876}
877
878pub(crate) enum RequestKind {
885 #[cfg(feature = "noq")]
886 Noq(Box<moq_net::Request<web_transport_noq::Session>>),
887 #[cfg(feature = "quinn")]
888 Quinn(Box<moq_net::Request<web_transport_quinn::Session>>),
889 #[cfg(feature = "quiche")]
890 Quiche(Box<moq_net::Request<web_transport_quiche::Connection>>),
891 #[cfg(feature = "iroh")]
892 Iroh(Box<moq_net::Request<web_transport_iroh::Session>>),
893 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
894 Qmux(Box<moq_net::Request<qmux::Session>>),
895}
896
897#[non_exhaustive]
899#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
900pub enum Transport {
901 Quic,
903 Iroh,
905 WebSocket,
907 Tcp,
909 Unix,
911}
912
913impl Transport {
914 pub const fn as_str(self) -> &'static str {
916 match self {
917 Self::Quic => "quic",
918 Self::Iroh => "iroh",
919 Self::WebSocket => "websocket",
920 Self::Tcp => "tcp",
921 Self::Unix => "unix",
922 }
923 }
924}
925
926impl std::fmt::Display for Transport {
927 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
928 f.write_str(self.as_str())
929 }
930}
931
932pub struct Request {
941 transport: Transport,
942 url: Option<Url>,
945 identity: Option<crate::tls::PeerIdentity>,
948 kind: RequestKind,
949}
950
951macro_rules! request_ref {
953 ($self:expr, $r:ident => $body:expr) => {
954 match &$self.kind {
955 #[cfg(feature = "noq")]
956 RequestKind::Noq($r) => $body,
957 #[cfg(feature = "quinn")]
958 RequestKind::Quinn($r) => $body,
959 #[cfg(feature = "quiche")]
960 RequestKind::Quiche($r) => $body,
961 #[cfg(feature = "iroh")]
962 RequestKind::Iroh($r) => $body,
963 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
964 RequestKind::Qmux($r) => $body,
965 }
966 };
967}
968
969macro_rules! request_into {
971 ($kind:expr, $r:ident => $body:expr) => {
972 match $kind {
973 #[cfg(feature = "noq")]
974 RequestKind::Noq($r) => $body,
975 #[cfg(feature = "quinn")]
976 RequestKind::Quinn($r) => $body,
977 #[cfg(feature = "quiche")]
978 RequestKind::Quiche($r) => $body,
979 #[cfg(feature = "iroh")]
980 RequestKind::Iroh($r) => $body,
981 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
982 RequestKind::Qmux($r) => $body,
983 }
984 };
985}
986
987macro_rules! request_map {
989 ($kind:expr, $r:ident => $body:expr) => {
990 match $kind {
991 #[cfg(feature = "noq")]
992 RequestKind::Noq($r) => RequestKind::Noq(Box::new($body)),
993 #[cfg(feature = "quinn")]
994 RequestKind::Quinn($r) => RequestKind::Quinn(Box::new($body)),
995 #[cfg(feature = "quiche")]
996 RequestKind::Quiche($r) => RequestKind::Quiche(Box::new($body)),
997 #[cfg(feature = "iroh")]
998 RequestKind::Iroh($r) => RequestKind::Iroh(Box::new($body)),
999 #[cfg(any(feature = "tcp", all(feature = "uds", unix), feature = "websocket"))]
1000 RequestKind::Qmux($r) => RequestKind::Qmux(Box::new($body)),
1001 }
1002 };
1003}
1004
1005impl Request {
1006 pub async fn close(self, code: u16) -> crate::Result<()> {
1010 let err = match code {
1011 401 | 403 => moq_net::Error::Unauthorized,
1012 other => moq_net::Error::App(other),
1013 };
1014 request_into!(self.kind, request => request.close(err));
1015 Ok(())
1016 }
1017
1018 pub fn with_publisher(self, publish: impl moq_net::Consume<moq_net::origin::Consumer>) -> Self {
1020 let Request {
1021 transport,
1022 url,
1023 identity,
1024 kind,
1025 } = self;
1026 let kind = request_map!(kind, request => request.with_publisher(publish));
1027 Request {
1028 transport,
1029 url,
1030 identity,
1031 kind,
1032 }
1033 }
1034
1035 pub fn with_subscriber(self, subscribe: moq_net::origin::Producer) -> Self {
1037 let Request {
1038 transport,
1039 url,
1040 identity,
1041 kind,
1042 } = self;
1043 let kind = request_map!(kind, request => request.with_subscriber(subscribe));
1044 Request {
1045 transport,
1046 url,
1047 identity,
1048 kind,
1049 }
1050 }
1051
1052 pub fn with_stats(self, stats: moq_net::stats::Session) -> Self {
1054 let Request {
1055 transport,
1056 url,
1057 identity,
1058 kind,
1059 } = self;
1060 let kind = request_map!(kind, request => request.with_stats(stats));
1061 Request {
1062 transport,
1063 url,
1064 identity,
1065 kind,
1066 }
1067 }
1068
1069 pub async fn ok(self) -> crate::Result<Session> {
1071 let pair = request_into!(self.kind, request => request.ok().await?);
1072 Ok(crate::spawn_session(pair))
1073 }
1074
1075 pub fn transport(&self) -> Transport {
1077 self.transport
1078 }
1079
1080 pub fn url(&self) -> Option<&Url> {
1085 self.url.as_ref()
1086 }
1087
1088 pub fn path(&self) -> &str {
1094 let setup = request_ref!(self, r => r.path());
1098 if setup.is_empty() {
1099 self.url.as_ref().map(Url::path).unwrap_or("")
1100 } else {
1101 setup
1102 }
1103 }
1104
1105 pub fn role(&self) -> Option<moq_net::Role> {
1110 request_ref!(self, r => r.role())
1111 }
1112
1113 pub fn peer_origin(&self) -> Option<moq_net::Origin> {
1121 request_ref!(self, r => r.peer_origin())
1122 }
1123
1124 pub fn peer_identity(&self) -> Option<crate::tls::PeerIdentity> {
1132 self.identity.clone()
1133 }
1134
1135 #[doc(hidden)]
1136 #[deprecated(note = "use `peer_identity` instead")]
1137 pub fn has_peer_certificate(&self) -> bool {
1138 self.peer_identity().is_some()
1139 }
1140}
1141
1142#[cfg(test)]
1143mod tests {
1144 use super::*;
1145
1146 #[cfg(feature = "tcp")]
1154 #[test]
1155 fn accept_health_covers_stream_listeners_before_they_bind() {
1156 let mut config = ServerConfig::default();
1157 config.tcp.bind = Some("127.0.0.1:0".parse().unwrap());
1158 let server = Server::new(config).expect("stream-only server");
1159
1160 let names: Vec<_> = server.accept_health().iter().map(|h| h.listener()).collect();
1161 assert_eq!(names, vec!["tcp"], "the tcp listener must report before it binds");
1162 }
1163
1164 #[cfg(all(feature = "tcp", feature = "uds", unix))]
1170 #[tokio::test]
1171 async fn a_failed_listen_binds_nothing_and_can_be_retried() {
1172 let dir = tempfile::TempDir::new().unwrap();
1175 let occupied = dir.path().join("not-a-socket");
1176 std::fs::write(&occupied, b"in the way").unwrap();
1177
1178 let mut config = ServerConfig::default();
1179 config.tcp.bind = Some("127.0.0.1:0".parse().unwrap());
1180 config.unix.bind = Some(occupied);
1181 let mut server = Server::new(config).expect("stream-only server");
1182
1183 assert!(server.listen().await.is_err(), "the unix bind must fail");
1184 assert!(server.listen().await.is_err(), "a retry must not report success");
1187 }
1188
1189 #[cfg(all(feature = "quinn", not(feature = "tcp")))]
1192 #[test]
1193 fn accept_health_is_empty_without_a_stream_listener() {
1194 let server = ServerConfig::default().init().expect("quic server");
1195 assert!(server.accept_health().is_empty());
1196 }
1197
1198 #[test]
1199 fn transport_names_are_stable() {
1200 assert_eq!(Transport::Quic.as_str(), "quic");
1201 assert_eq!(Transport::Iroh.as_str(), "iroh");
1202 assert_eq!(Transport::WebSocket.as_str(), "websocket");
1203 assert_eq!(Transport::Tcp.as_str(), "tcp");
1204 assert_eq!(Transport::Unix.as_str(), "unix");
1205 }
1206
1207 #[cfg(feature = "quinn")]
1210 #[tokio::test]
1211 async fn certificates_expose_generated_fingerprints() {
1212 let mut config = ServerConfig {
1213 bind: Some("[::]:0".to_string()),
1214 ..Default::default()
1215 };
1216 config.tls.generate = vec!["localhost".into()];
1217
1218 let certs = config.init().expect("server init").certificates();
1219 let fingerprints = certs.fingerprints();
1220 assert_eq!(fingerprints.len(), 1, "one generated certificate");
1221 assert_eq!(fingerprints[0].len(), 64);
1223 assert!(fingerprints[0].chars().all(|c| c.is_ascii_hexdigit()));
1224 }
1225
1226 #[cfg(all(feature = "uds", unix))]
1231 #[tokio::test]
1232 async fn unix_listener_serves_the_configured_publisher() {
1233 use rand::RngExt;
1234
1235 let path = PathBuf::from(format!("/tmp/moq-native-publish-{}.sock", std::process::id()));
1238 let _ = std::fs::remove_file(&path);
1239
1240 let origin = moq_net::Origin::random().produce();
1241 let mut broadcast = origin
1242 .create_broadcast("test", moq_net::broadcast::Route::new().with_announce(true))
1243 .expect("create broadcast");
1244 let mut track = broadcast.create_track("video", None).expect("create track");
1245 let mut group = track.append_group().expect("append group");
1246 group
1247 .write_frame(moq_net::Timestamp::ZERO, b"hello".as_ref())
1248 .expect("write frame");
1249 group.finish().expect("finish group");
1250
1251 let mut config = ServerConfig::default();
1252 config.unix.bind = Some(path.clone());
1253 let server = config.init().expect("server init");
1254
1255 let serve = tokio::spawn(server.serve_publish(origin.consume()));
1257
1258 const MAX_DELAY: std::time::Duration = std::time::Duration::from_millis(100);
1262 let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
1263 let mut delay = std::time::Duration::from_millis(1);
1264 while let Err(err) = tokio::net::UnixStream::connect(&path).await {
1265 assert!(
1266 tokio::time::Instant::now() < deadline,
1267 "unix listener never bound: {err}"
1268 );
1269 tokio::time::sleep(delay.mul_f64(0.5 + rand::rng().random::<f64>() / 2.0)).await;
1270 delay = (delay * 2).min(MAX_DELAY);
1271 }
1272
1273 const TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
1274
1275 let url: Url = format!("unix://{}", path.display()).parse().expect("parse url");
1276 let subscriber = moq_net::Origin::random().produce();
1277 let mut announced = subscriber.consume().announced();
1278 let client = crate::ClientConfig::default()
1279 .init()
1280 .expect("client init")
1281 .with_subscriber(subscriber);
1282 let session = tokio::time::timeout(TIMEOUT, client.connect(url))
1283 .await
1284 .expect("connect timeout")
1285 .expect("connect");
1286
1287 let update = tokio::time::timeout(TIMEOUT, announced.next())
1290 .await
1291 .expect("announce timeout")
1292 .expect("origin closed");
1293 assert_eq!(update.path.as_str(), "test");
1294 let broadcast = update.broadcast.expect("expected an announce");
1295
1296 let mut track = broadcast
1297 .track("video")
1298 .expect("track name")
1299 .subscribe(None)
1300 .await
1301 .expect("subscribe");
1302 let mut group = tokio::time::timeout(TIMEOUT, track.recv_group())
1303 .await
1304 .expect("recv group timeout")
1305 .expect("recv group")
1306 .expect("track closed early");
1307 let frame = tokio::time::timeout(TIMEOUT, group.read_frame())
1308 .await
1309 .expect("read frame timeout")
1310 .expect("read frame")
1311 .expect("group closed early");
1312 assert_eq!(&frame.payload[..], b"hello");
1313
1314 drop(session);
1315 serve.abort();
1316 let _ = std::fs::remove_file(&path);
1317 }
1318
1319 #[cfg(all(feature = "uds", unix))]
1322 #[tokio::test]
1323 async fn certificates_are_empty_without_a_tls_backend() {
1324 let mut config = ServerConfig::default();
1325 config.unix.bind = Some(PathBuf::from("/tmp/moq-native-certificates-test.sock"));
1326
1327 let server = config.init().expect("server init");
1328 assert!(server.certificates().fingerprints().is_empty());
1329 }
1330
1331 #[test]
1332 fn test_tls_string_or_array() {
1333 let single = r#"
1335 cert = "cert.pem"
1336 key = "key.pem"
1337 "#;
1338 let config: crate::tls::Server = toml::from_str(single).unwrap();
1339 assert_eq!(config.cert, vec![PathBuf::from("cert.pem")]);
1340 assert_eq!(config.key, vec![PathBuf::from("key.pem")]);
1341
1342 let array = r#"
1344 cert = ["a.pem", "b.pem"]
1345 key = ["a.key", "b.key"]
1346 generate = ["localhost"]
1347 root = ["ca.pem"]
1348 "#;
1349 let config: crate::tls::Server = toml::from_str(array).unwrap();
1350 assert_eq!(config.cert, vec![PathBuf::from("a.pem"), PathBuf::from("b.pem")]);
1351 assert_eq!(config.key, vec![PathBuf::from("a.key"), PathBuf::from("b.key")]);
1352 assert_eq!(config.generate, vec!["localhost".to_string()]);
1353 assert_eq!(config.root, vec![PathBuf::from("ca.pem")]);
1354 }
1355
1356 #[test]
1357 fn bind_string_or_listen_alias() {
1358 let bind: ServerConfig = toml::from_str(r#"bind = "[::]:443""#).unwrap();
1360 assert_eq!(bind.bind.as_deref(), Some("[::]:443"));
1361
1362 let alias: ServerConfig = toml::from_str(r#"listen = "0.0.0.0:4443""#).unwrap();
1363 assert_eq!(alias.bind.as_deref(), Some("0.0.0.0:4443"));
1364 }
1365
1366 #[cfg(all(feature = "uds", unix))]
1367 #[test]
1368 fn stream_listener_config_parses() {
1369 let config: ServerConfig = toml::from_str(
1370 r#"
1371bind = "[::]:443"
1372
1373[unix]
1374bind = "/run/moq.sock"
1375
1376[unix.allow]
1377uid = [1001, 1002]
1378"#,
1379 )
1380 .unwrap();
1381 assert_eq!(config.bind.as_deref(), Some("[::]:443"));
1382 assert_eq!(config.unix.bind.as_deref(), Some(std::path::Path::new("/run/moq.sock")));
1383 assert_eq!(config.unix.allow.as_ref().expect("allow").uid, vec![1001, 1002]);
1384 assert!(config.has_stream_listener());
1385 }
1386
1387 #[cfg(all(feature = "uds", unix))]
1388 #[test]
1389 fn stream_only_config_has_no_quic() {
1390 let mut config = ServerConfig::default();
1392 config.unix.bind = Some(PathBuf::from("/run/moq.sock"));
1393 assert!(config.has_stream_listener());
1394 assert!(config.bind.is_none());
1395
1396 assert!(!ServerConfig::default().has_stream_listener());
1398 }
1399}