use std::task::Poll;
use serde::de::DeserializeOwned;
use super::decoder::Codec;
use super::{ConsumerConfig, Decoder, Event};
use crate::Result;
pub struct Consumer<T> {
track: moq_net::track::Subscriber,
group: Option<moq_net::group::Consumer>,
codec: Option<Codec>,
decoder: Decoder<T>,
}
impl<T: DeserializeOwned> Consumer<T> {
pub fn new(track: moq_net::track::Subscriber, config: ConsumerConfig) -> Self {
Self {
track,
group: None,
codec: None,
decoder: Decoder::new(config),
}
}
pub fn range(&self) -> std::ops::Range<u64> {
self.decoder.range()
}
pub async fn next(&mut self) -> Result<Option<Event<T>>>
where
T: Unpin,
{
kio::wait(|waiter| self.poll_next(waiter)).await
}
pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Event<T>>>> {
loop {
if let Some(event) = self.decoder.next_event() {
return Poll::Ready(Ok(Some(event)));
}
let Some(group) = &mut self.group else {
match self.track.poll_next_group(waiter)? {
Poll::Ready(Some(group)) => {
self.codec = Some(Codec::new());
self.group = Some(group);
continue;
}
Poll::Ready(None) => return Poll::Ready(Ok(None)),
Poll::Pending => return Poll::Pending,
}
};
match group.poll_read_frame(waiter)? {
Poll::Ready(Some(frame)) => {
let codec = self.codec.as_mut().expect("an open MoQ group has a window codec");
self.decoder.decode(codec, &frame.payload)?;
}
Poll::Ready(None) => {
self.group = None;
self.codec = None;
}
Poll::Pending => return Poll::Pending,
}
}
}
}