moq_json/window/
consumer.rs1use std::task::Poll;
4
5use serde::de::DeserializeOwned;
6
7use super::decoder::Codec;
8use super::{ConsumerConfig, Decoder, Event};
9use crate::Result;
10
11pub struct Consumer<T> {
21 track: moq_net::track::Subscriber,
22 group: Option<moq_net::group::Consumer>,
23 codec: Option<Codec>,
24 decoder: Decoder<T>,
25}
26
27impl<T: DeserializeOwned> Consumer<T> {
28 pub fn new(track: moq_net::track::Subscriber, config: ConsumerConfig) -> Self {
30 Self {
31 track,
32 group: None,
33 codec: None,
34 decoder: Decoder::new(config),
35 }
36 }
37
38 pub fn range(&self) -> std::ops::Range<u64> {
40 self.decoder.range()
41 }
42
43 pub async fn next(&mut self) -> Result<Option<Event<T>>>
45 where
46 T: Unpin,
47 {
48 kio::wait(|waiter| self.poll_next(waiter)).await
49 }
50
51 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Event<T>>>> {
53 loop {
54 if let Some(event) = self.decoder.next_event() {
56 return Poll::Ready(Ok(Some(event)));
57 }
58
59 let Some(group) = &mut self.group else {
60 match self.track.poll_next_group(waiter)? {
61 Poll::Ready(Some(group)) => {
62 self.codec = Some(Codec::new());
63 self.group = Some(group);
64 continue;
65 }
66 Poll::Ready(None) => return Poll::Ready(Ok(None)),
67 Poll::Pending => return Poll::Pending,
68 }
69 };
70
71 match group.poll_read_frame(waiter)? {
72 Poll::Ready(Some(frame)) => {
73 let codec = self.codec.as_mut().expect("an open MoQ group has a window codec");
74 self.decoder.decode(codec, &frame.payload)?;
75 }
76 Poll::Ready(None) => {
77 self.group = None;
80 self.codec = None;
81 }
82 Poll::Pending => return Poll::Pending,
83 }
84 }
85 }
86}