Skip to main content

moq_rtc/
lib.rs

1//! WebRTC ↔ MoQ gateway.
2//!
3//! Bridges WHIP (RFC 9725) and WHEP between WebRTC peers and
4//! [`moq_net`] broadcasts. The crate is split along two orthogonal axes
5//! so all four combinations can land independently:
6//!
7//! | | RTP-in (ingest into MoQ) | RTP-out (egress from MoQ) |
8//! |---|---|---|
9//! | HTTP server | [`Server::publish_router`] (WHIP server) | [`Server::subscribe_router`] (WHEP server) |
10//! | HTTP client | [`Client::subscribe`] (WHEP client) | [`Client::publish`] (WHIP client) |
11//!
12//! The two HTTP-client paths and the two HTTP-server paths share a single
13//! internal session driver and the same per-codec adapters; the per-direction
14//! split lives in the (crate-private) ingest and egress sources.
15//!
16//! ## Embedding
17//!
18//! Build a [`Server`] and pass your own origin handles when merging
19//! [`Server::publish_router`] / [`Server::subscribe_router`] into your own axum
20//! app, or dial out with [`Client`]. A command-line interface is provided by the
21//! `moq-cli` binary, on top of this library.
22//!
23//! The bundled routers are unauthenticated: they derive the broadcast name from
24//! the request path. To own the HTTP route and authorize requests yourself
25//! (resolving the broadcast name from a verified token), skip the routers and
26//! call [`whip::accept`] (ingest) / [`whep::accept`] (egress) from your own
27//! handler. Return the [`Response::answer`] in your HTTP response, then run
28//! [`Response::run`] to drive the media session for its lifetime. The routers and
29//! the `axum` re-export sit behind the default `server` feature, so such an
30//! embedder can drop them with `default-features = false`.
31//!
32//! ## Bitstream gotcha
33//!
34//! The WebRTC ↔ MoQ shape conversion for H.264 and H.265 is handled by
35//! `moq-mux` importers: str0m hands us Annex-B (start-code NALs with inline
36//! parameter sets) and that's exactly what the importers want. AV1 uses the
37//! shared OBU splitter/importer. Opus, VP8, and VP9 pass through.
38
39#![warn(missing_docs)]
40
41pub mod client;
42pub mod server;
43
44// Implementation detail modules: these carry the WebRTC/str0m plumbing (str0m
45// `Rtc`, `Mid`/`Pt`, tokio channels, raw packet buffers) and are deliberately
46// crate-private, so the public surface stays `Client`, `Server`,
47// `whip`/`whep::accept`, and `Response`.
48mod codec;
49mod egress;
50mod error;
51mod ingest;
52mod net;
53mod sdp;
54mod session;
55
56/// Re-export of the HTTP router stack, so consumers can merge the [`axum::Router`]
57/// returned by [`Server::publish_router`] / [`Server::subscribe_router`] (and by
58/// [`whip::router`] / [`whep::router`]) into their own app without adding their own
59/// axum dependency (and risking a version mismatch). A major axum bump is therefore
60/// a breaking change for this crate. Only with the `server` feature.
61#[cfg(feature = "server")]
62pub use axum;
63
64/// Re-export of the URL type, so consumers can build the [`url::Url`] that
65/// [`Client::subscribe`] / [`Client::publish`] dial without adding their own url
66/// dependency (and risking a version mismatch). A major url bump is therefore a
67/// breaking change for this crate.
68pub use url;
69
70pub use client::Client;
71pub use error::*;
72pub use server::{Response, Server, whep, whip};
73
74#[cfg(all(test, feature = "server"))]
75mod tests {
76	use std::time::Duration;
77
78	use axum::Router;
79	use bytes::Bytes;
80
81	use crate::codec::{Bridge, Frame, Track};
82	use crate::{Client, Server, client, server};
83
84	const TIMEOUT: Duration = Duration::from_secs(10);
85	const OPUS_PACKET: &[u8] = &[0xfc, 0xff, 0xfe];
86
87	#[tokio::test]
88	#[tracing_test::traced_test]
89	async fn whip_and_whep_round_trip_opus() {
90		let source_origin = moq_tokio::origin::spawn();
91		let source_consumer = source_origin.consume();
92		let mut announcements = source_consumer.announced();
93		let mut source = source_origin
94			.create_broadcast("source")
95			.expect("create source broadcast");
96		source
97			.announce(moq_net::origin::Route::default())
98			.expect("announce source broadcast");
99		let catalog = moq_mux::catalog::Producer::new(&mut source, moq_mux::catalog::Config::default())
100			.expect("create source catalog");
101		let mut opus = crate::codec::opus::Bridge::new(source, catalog, 48_000, 2).expect("create Opus bridge");
102		Bridge::push(
103			&mut opus,
104			Frame {
105				timestamp_us: 20_000,
106				payload: Bytes::from_static(OPUS_PACKET),
107			},
108		)
109		.expect("publish source packet");
110		let announcement = tokio::time::timeout(TIMEOUT, announcements.next())
111			.await
112			.expect("source announcement timed out")
113			.expect("source origin closed");
114		assert_eq!(announcement.prefix.as_str(), "source");
115		assert!(announcement.kind.is_active(), "source was unannounced");
116		drop(announcements);
117
118		let server_origin = moq_tokio::origin::spawn();
119		let server = Server::new(server::Config::default());
120		let app = Router::new()
121			.nest("/whip", server.publish_router(server_origin.clone()))
122			.nest("/whep", server.subscribe_router(server_origin.consume()));
123		let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
124			.await
125			.expect("bind HTTP listener");
126		let address = listener.local_addr().expect("HTTP listener address");
127		let http = tokio::spawn(async move { axum::serve(listener, app).await.expect("serve HTTP") });
128
129		let client = Client::new(client::Config::default());
130		// Credentials in the dialed URLs must never reach the client's connect logs.
131		let whip = format!("http://user:pass@{address}/whip/ingested?jwt=secret")
132			.parse()
133			.expect("WHIP URL");
134		tokio::time::timeout(TIMEOUT, client.publish(whip, source_consumer, "source"))
135			.await
136			.expect("WHIP negotiation timed out")
137			.expect("WHIP negotiation failed");
138
139		let output_origin = moq_tokio::origin::spawn();
140		let output = output_origin
141			.create_broadcast("output")
142			.expect("create output broadcast");
143		let output_consumer = output.consume();
144		let whep = format!("http://user:pass@{address}/whep/ingested?jwt=secret")
145			.parse()
146			.expect("WHEP URL");
147		tokio::time::timeout(TIMEOUT, client.subscribe(whep, output))
148			.await
149			.expect("WHEP negotiation timed out")
150			.expect("WHEP negotiation failed");
151
152		assert!(logs_contain("whip client connected"));
153		assert!(logs_contain("whep client connected"));
154		for secret in ["jwt", "secret", "user:pass"] {
155			assert!(!logs_contain(secret), "client logs leaked {secret}");
156		}
157
158		let catalog_track = output_consumer
159			.track(hang::Catalog::DEFAULT_NAME)
160			.expect("output catalog track")
161			.subscribe(hang::Catalog::default_subscription())
162			.await
163			.expect("subscribe to output catalog");
164		let mut catalogs = moq_mux::catalog::hang::Consumer::<()>::new(catalog_track);
165		let catalog = tokio::time::timeout(TIMEOUT, catalogs.next())
166			.await
167			.expect("output catalog timed out")
168			.expect("read output catalog")
169			.expect("output catalog ended");
170		let track_name = catalog.audio.renditions.keys().next().expect("output Opus rendition");
171		let track = output_consumer
172			.track(track_name)
173			.expect("output Opus track")
174			.subscribe(None)
175			.await
176			.expect("subscribe to output Opus track");
177		let mut opus = Track::opus(track);
178		let frame = tokio::time::timeout(TIMEOUT, opus.next())
179			.await
180			.expect("output Opus packet timed out")
181			.expect("read output Opus packet")
182			.expect("output Opus track ended");
183		assert_eq!(frame.payload.as_ref(), OPUS_PACKET);
184
185		http.abort();
186	}
187}