moq_json/stream/
producer.rs1use std::marker::PhantomData;
4use std::sync::{Arc, Mutex};
5
6use serde::Serialize;
7
8use super::{Encoder, ProducerConfig};
9use crate::Result;
10
11pub 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 consume(&self) -> moq_net::track::Subscriber {
35 self.inner.lock().unwrap().track.inner.subscribe(None)
36 }
37}
38
39impl<T: Serialize> Producer<T> {
40 pub fn new(track: moq_net::track::Producer, config: ProducerConfig) -> Self {
42 Self {
43 inner: Arc::new(Mutex::new(Inner {
44 track: Track {
45 inner: track,
46 group: None,
47 },
48 encoder: Encoder::new(config),
49 })),
50 _marker: PhantomData,
51 }
52 }
53
54 pub fn append(&mut self, value: &T) -> Result<()> {
56 self.inner.lock().unwrap().append(value)
57 }
58
59 pub fn finish(&mut self) -> Result<()> {
61 self.inner.lock().unwrap().finish()
62 }
63}
64
65struct Inner<T> {
71 track: Track,
72 encoder: Encoder<T>,
73}
74
75impl<T: Serialize> Inner<T> {
76 fn append(&mut self, value: &T) -> Result<()> {
77 let Inner { track, encoder } = self;
79
80 let record = encoder.encode(value)?;
84
85 let result = match track.open() {
86 Ok(()) => track.write(record.payload()),
87 Err(err) => Err(err),
88 };
89
90 if let Err(err) = result {
91 drop(record);
96 encoder.reset();
97 return Err(err);
98 }
99
100 record.commit();
101 Ok(())
102 }
103
104 fn finish(&mut self) -> Result<()> {
105 self.track.finish()
106 }
107}
108
109struct Track {
111 inner: moq_net::track::Producer,
112 group: Option<moq_net::group::Producer>,
114}
115
116impl Track {
117 fn open(&mut self) -> Result<()> {
119 if self.group.is_none() {
120 self.group = Some(self.inner.append_group()?);
121 }
122 Ok(())
123 }
124
125 fn write(&mut self, payload: &bytes::Bytes) -> Result<()> {
127 let group = self.group.as_mut().expect("a group is open");
128 let Err(err) = group.write_frame(moq_net::Timestamp::now(), payload.clone()) else {
129 return Ok(());
130 };
131
132 if let Some(mut group) = self.group.take() {
136 let _ = group.finish();
137 }
138 Err(err.into())
139 }
140
141 fn finish(&mut self) -> Result<()> {
142 if let Some(mut group) = self.group.take() {
143 group.finish()?;
144 }
145 self.inner.finish()?;
146 Ok(())
147 }
148}