use super::*;
use crate::error::flatten_errs;
use tokio::sync::Mutex;
#[derive(Debug)]
pub struct TrackLocalStaticRTP {
pub(crate) bindings: Mutex<Vec<Arc<TrackBinding>>>,
codec: RTCRtpCodecCapability,
id: String,
stream_id: String,
}
impl TrackLocalStaticRTP {
pub fn new(codec: RTCRtpCodecCapability, id: String, stream_id: String) -> Self {
TrackLocalStaticRTP {
codec,
bindings: Mutex::new(vec![]),
id,
stream_id,
}
}
pub fn codec(&self) -> RTCRtpCodecCapability {
self.codec.clone()
}
}
#[async_trait]
impl TrackLocal for TrackLocalStaticRTP {
async fn bind(&self, t: &TrackLocalContext) -> Result<RTCRtpCodecParameters> {
let parameters = RTCRtpCodecParameters {
capability: self.codec.clone(),
..Default::default()
};
let (codec, match_type) = codec_parameters_fuzzy_search(¶meters, t.codec_parameters());
if match_type != CodecMatch::None {
{
let mut bindings = self.bindings.lock().await;
bindings.push(Arc::new(TrackBinding {
ssrc: t.ssrc(),
payload_type: codec.payload_type,
write_stream: t.write_stream(),
id: t.id(),
}));
}
Ok(codec)
} else {
Err(Error::ErrUnsupportedCodec)
}
}
async fn unbind(&self, t: &TrackLocalContext) -> Result<()> {
let mut bindings = self.bindings.lock().await;
let mut idx = None;
for (index, binding) in bindings.iter().enumerate() {
if binding.id == t.id() {
idx = Some(index);
break;
}
}
if let Some(index) = idx {
bindings.remove(index);
Ok(())
} else {
Err(Error::ErrUnbindFailed)
}
}
fn id(&self) -> &str {
self.id.as_str()
}
fn stream_id(&self) -> &str {
self.stream_id.as_str()
}
fn kind(&self) -> RTPCodecType {
if self.codec.mime_type.starts_with("audio/") {
RTPCodecType::Audio
} else if self.codec.mime_type.starts_with("video/") {
RTPCodecType::Video
} else {
RTPCodecType::Unspecified
}
}
fn as_any(&self) -> &dyn Any {
self
}
}
#[async_trait]
impl TrackLocalWriter for TrackLocalStaticRTP {
async fn write_rtp(&self, p: &rtp::packet::Packet) -> Result<usize> {
let mut n = 0;
let mut write_errs = vec![];
let mut pkt = p.clone();
let bindings = {
let bindings = self.bindings.lock().await;
bindings.clone()
};
for b in bindings {
pkt.header.ssrc = b.ssrc;
pkt.header.payload_type = b.payload_type;
if let Some(write_stream) = &b.write_stream {
match write_stream.write_rtp(&pkt).await {
Ok(m) => {
n += m;
}
Err(err) => {
write_errs.push(err);
}
}
} else {
write_errs.push(Error::new("track binding has none write_stream".to_owned()));
}
}
flatten_errs(write_errs)?;
Ok(n)
}
async fn write(&self, mut b: &[u8]) -> Result<usize> {
let pkt = rtp::packet::Packet::unmarshal(&mut b)?;
self.write_rtp(&pkt).await?;
Ok(b.len())
}
}