use std::time::Instant;
use str0m::{
Candidate, Rtc,
change::SdpAnswer,
media::{Direction, MediaKind},
};
use url::Url;
use crate::{Error, Result, client::Client, ingest::IngestSink, session};
pub(crate) async fn dial(client: &Client, url: Url, broadcast: moq_net::broadcast::Producer) -> Result<()> {
let config = moq_mux::catalog::Config::default()
.with_max_age(client.config().max_age)
.with_bandwidth(client.config().bandwidth.clone());
let sink = Box::new(IngestSink::new(broadcast, config)?);
let (socket, candidates) = session::bind_udp(&client.config().ice_candidates).await?;
let mut rtc = Rtc::new(Instant::now());
for addr in &candidates {
let cand = Candidate::host(*addr, "udp").map_err(Error::rtc)?;
rtc.add_local_candidate(cand);
}
let mut api = rtc.sdp_api();
let mids = [
api.add_media(MediaKind::Audio, Direction::RecvOnly, None, None, None),
api.add_media(MediaKind::Video, Direction::RecvOnly, None, None, None),
];
let (offer, pending) = api.apply().ok_or(Error::NoSdpChanges)?;
let res = client
.http()
.post(url.clone())
.header(reqwest::header::CONTENT_TYPE, "application/sdp")
.header(reqwest::header::ACCEPT, "application/sdp")
.body(offer.to_sdp_string())
.send()
.await?;
if !res.status().is_success() {
return Err(Error::HttpStatus(res.status().as_u16()));
}
let body = res.text().await?;
let answer = SdpAnswer::from_sdp_string(&body).map_err(|err| Error::InvalidSdp(err.to_string()))?;
rtc.sdp_api().accept_answer(pending, answer).map_err(Error::rtc)?;
tracing::info!(%url, "whep client connected");
let inbound = session::spawn_socket_reader(socket.clone());
let mut session = session::Session::ingest(rtc, socket, candidates, inbound, sink);
for mid in mids {
session.open_media(mid)?;
}
tokio::spawn(async move {
let result = session.run().await;
session::log_session_end("whep client", &result);
});
Ok(())
}