1#![warn(missing_docs)]
40
41pub mod client;
42pub mod server;
43
44mod codec;
49mod egress;
50mod error;
51mod ingest;
52mod net;
53mod sdp;
54mod session;
55
56#[cfg(feature = "server")]
62pub use axum;
63
64pub 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 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}