use super::decoder::FlexFec03Decoder;
use crate::Interceptor;
use crate::stream_info::StreamInfo;
use crate::{Attribute, AttributedPacket, Packet, TaggedPacket};
use sansio::Protocol;
use shared::error::Error;
use std::collections::{HashMap, VecDeque};
use std::time::Instant;
#[derive(Default)]
pub struct FlexFec03ReceiveBuilder {}
impl FlexFec03ReceiveBuilder {
pub fn new() -> Self {
Self::default()
}
pub fn build(self) -> FlexFec03ReceiveInterceptor {
FlexFec03ReceiveInterceptor::new()
}
}
pub struct FlexFec03ReceiveInterceptor {
decoders: HashMap<u32, FlexFec03Decoder>,
repair_to_media: HashMap<u32, u32>,
read_queue: VecDeque<TaggedPacket>,
write_queue: VecDeque<TaggedPacket>,
}
impl FlexFec03ReceiveInterceptor {
fn new() -> Self {
Self {
read_queue: VecDeque::new(),
write_queue: VecDeque::new(),
decoders: HashMap::new(),
repair_to_media: HashMap::new(),
}
}
pub fn protected_streams(&self) -> impl Iterator<Item = u32> + '_ {
self.decoders.keys().copied()
}
}
impl Protocol<TaggedPacket, TaggedPacket, ()> for FlexFec03ReceiveInterceptor {
type Rout = TaggedPacket;
type Wout = TaggedPacket;
type Eout = ();
type Error = Error;
type Time = Instant;
fn handle_read(&mut self, msg: TaggedPacket) -> Result<(), Self::Error> {
let Packet::Rtp(rtp_packet) = &msg.message.packet else {
self.read_queue.push_back(msg);
return Ok(());
};
let ssrc = rtp_packet.header.ssrc;
let now = msg.now;
let transport = msg.transport;
if let Some(&media_ssrc) = self.repair_to_media.get(&ssrc) {
let recovered = match self.decoders.get_mut(&media_ssrc) {
Some(decoder) => decoder.decode(rtp_packet.clone()),
None => Vec::new(),
};
self.queue_recovered(now, transport, recovered);
return Ok(());
}
let recovered = match self.decoders.get_mut(&ssrc) {
Some(decoder) => decoder.decode(rtp_packet.clone()),
None => {
self.read_queue.push_back(msg);
return Ok(());
}
};
self.read_queue.push_back(msg);
self.queue_recovered(now, transport, recovered);
Ok(())
}
fn poll_read(&mut self) -> Option<TaggedPacket> {
self.read_queue.pop_front()
}
fn handle_write(&mut self, msg: TaggedPacket) -> Result<(), Self::Error> {
self.write_queue.push_back(msg);
Ok(())
}
fn poll_write(&mut self) -> Option<Self::Wout> {
self.write_queue.pop_front()
}
fn handle_timeout(&mut self, _now: Instant) -> Result<(), Self::Error> {
Ok(())
}
fn poll_timeout(&mut self) -> Option<Self::Time> {
None
}
}
impl Interceptor for FlexFec03ReceiveInterceptor {
fn bind_remote_stream(&mut self, info: &StreamInfo) {
if let (Some(ssrc_fec), Some(_)) = (info.ssrc_fec, info.payload_type_fec) {
self.decoders
.insert(info.ssrc, FlexFec03Decoder::new(ssrc_fec, info.ssrc));
self.repair_to_media.insert(ssrc_fec, info.ssrc);
}
}
fn unbind_remote_stream(&mut self, info: &StreamInfo) {
self.decoders.remove(&info.ssrc);
self.repair_to_media.retain(|_, media| *media != info.ssrc);
}
fn bind_local_stream(&mut self, _info: &StreamInfo) {}
fn unbind_local_stream(&mut self, _info: &StreamInfo) {}
}
impl FlexFec03ReceiveInterceptor {
fn queue_recovered(
&mut self,
now: std::time::Instant,
transport: shared::TransportContext,
recovered: Vec<rtp::Packet>,
) {
for packet in recovered {
let mut message = AttributedPacket::new(Packet::Rtp(packet));
message.add(Attribute::RecoveredByFec);
self.read_queue.push_back(TaggedPacket {
now,
transport,
message,
});
}
}
}