moq_json/window/
consumer.rs1use std::task::{Poll, ready};
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 ready!(self.track.poll_next_group(waiter)?) {
61 Some(group) => {
62 self.codec = Some(Codec::new());
63 self.group = Some(group);
64 continue;
65 }
66 None => return Poll::Ready(Ok(None)),
67 }
68 };
69
70 match ready!(group.poll_read_frame(waiter)?) {
71 Some(frame) => {
72 let codec = self.codec.as_mut().expect("an open MoQ group has a window codec");
73 self.decoder.decode(codec, &frame.payload)?;
74 }
75 None => {
76 self.group = None;
79 self.codec = None;
80 }
81 }
82 }
83 }
84}