moq_json/window/
producer.rs1use std::marker::PhantomData;
4use std::sync::{Arc, Mutex};
5
6use serde::Serialize;
7use serde_json::Value;
8
9use super::{Encoded, Encoder, ProducerConfig};
10use crate::Result;
11
12pub struct Producer<T> {
19 inner: Arc<Mutex<Inner<T>>>,
20 _marker: PhantomData<fn(T)>,
21}
22
23impl<T> Clone for Producer<T> {
24 fn clone(&self) -> Self {
25 Self {
26 inner: self.inner.clone(),
27 _marker: PhantomData,
28 }
29 }
30}
31
32impl<T> Producer<T> {
33 pub fn new(track: moq_net::track::Producer, config: ProducerConfig) -> Self {
35 Self {
36 inner: Arc::new(Mutex::new(Inner {
37 track: Track {
38 inner: track,
39 group: None,
40 },
41 encoder: Encoder::new(config),
42 finished: false,
43 })),
44 _marker: PhantomData,
45 }
46 }
47
48 pub fn consume(&self) -> moq_net::track::Subscriber {
50 self.inner.lock().unwrap().track.inner.subscribe(None)
51 }
52
53 pub fn window(&self) -> Vec<Value> {
57 self.inner.lock().unwrap().encoder.window()
58 }
59
60 pub fn range(&self) -> std::ops::Range<u64> {
62 self.inner.lock().unwrap().encoder.range()
63 }
64
65 pub fn pop(&mut self, count: u64) -> Result<()> {
70 self.inner.lock().unwrap().pop(count)
71 }
72
73 pub fn finish(&mut self) -> Result<()> {
79 self.inner.lock().unwrap().finish()
80 }
81}
82
83impl<T: Serialize> Producer<T> {
84 pub fn push(&mut self, value: &T) -> Result<()> {
86 self.inner.lock().unwrap().push(value)
87 }
88}
89
90struct Inner<T> {
96 track: Track,
97 encoder: Encoder<T>,
98 finished: bool,
99}
100
101impl<T> Inner<T> {
102 fn pop(&mut self, count: u64) -> Result<()> {
103 self.ensure_open()?;
104 let Inner { track, encoder, .. } = self;
105
106 let Some(frame) = encoder.pop(count)? else {
107 return Ok(());
108 };
109
110 track.write(&frame)?;
113 frame.commit();
114
115 Ok(())
116 }
117
118 fn finish(&mut self) -> Result<()> {
119 if self.finished {
120 return Ok(());
121 }
122 self.finished = true;
123 self.track.finish()
124 }
125
126 fn ensure_open(&self) -> Result<()> {
127 if self.finished {
128 return Err(moq_net::Error::Closed.into());
129 }
130 Ok(())
131 }
132}
133
134impl<T: Serialize> Inner<T> {
135 fn push(&mut self, value: &T) -> Result<()> {
136 self.ensure_open()?;
137 let Inner { track, encoder, .. } = self;
138
139 let frame = encoder.push(value)?;
140 track.write(&frame)?;
141 frame.commit();
142
143 Ok(())
144 }
145}
146
147struct Track {
149 inner: moq_net::track::Producer,
150
151 group: Option<moq_net::group::Producer>,
153}
154
155impl Track {
156 fn write(&mut self, encoded: &Encoded) -> Result<()> {
158 match encoded.keyframe {
159 true => self.write_header(encoded.payload.clone()),
160 false => self.write_op(encoded.payload.clone()),
161 }
162 }
163
164 fn write_header(&mut self, payload: bytes::Bytes) -> Result<()> {
166 if let Some(group) = self.group.take() {
168 group.finish()?;
169 }
170
171 let mut group = self.inner.append_group()?;
172 if let Err(err) = group.write_frame(moq_net::Timestamp::now(), payload) {
173 let _ = group.finish();
177 return Err(err.into());
178 }
179
180 self.group = Some(group);
181 Ok(())
182 }
183
184 fn write_op(&mut self, payload: bytes::Bytes) -> Result<()> {
186 self.group
187 .as_mut()
188 .expect("the encoder only emits an op after a header opened a group")
189 .write_frame(moq_net::Timestamp::now(), payload)?;
190 Ok(())
191 }
192
193 fn finish(&mut self) -> Result<()> {
194 if let Some(group) = self.group.take() {
195 group.finish()?;
196 }
197 self.inner.finish()?;
198 Ok(())
199 }
200}