moq_binary/stream/producer.rs
1//! Publishing an ordered log of binary payloads over a track.
2
3use std::sync::{Arc, Mutex};
4
5use bytes::Bytes;
6use moq_net::Timed;
7
8use crate::Result;
9
10pub use super::Config;
11
12/// Publishes an ordered log of binary payloads over a track, one payload per frame in a single
13/// group.
14///
15/// Cheaply clonable: clones share one underlying track and publishing state, so multiple owners
16/// append into a single ordered log.
17#[derive(Clone)]
18pub struct Producer {
19 inner: Arc<Mutex<Inner>>,
20}
21
22impl Producer {
23 /// Create a producer that publishes to the given track.
24 pub fn new(track: moq_net::track::Producer, config: Config) -> Self {
25 Self {
26 inner: Arc::new(Mutex::new(Inner {
27 track,
28 group: None,
29 flate: config.compression.is_deflate().then(moq_flate::Encoder::new),
30 })),
31 }
32 }
33
34 /// Create a subscriber for the underlying track.
35 ///
36 /// Still hands one back once a failed write has ended the log: the subscriber surfaces the abort
37 /// on its first read, which is what tells a late reader the log is truncated.
38 pub fn consume(&self) -> moq_net::track::Subscriber {
39 self.inner.lock().unwrap().track.subscribe(None)
40 }
41
42 /// Whether any consumer for the underlying track currently exists.
43 ///
44 /// The demand signal for a producer serving on request: an unused track is cached state nobody is
45 /// watching, safe to drop and recreate on the next request.
46 pub fn is_used(&self) -> bool {
47 self.inner.lock().unwrap().track.is_used()
48 }
49
50 /// Append one payload to the log.
51 ///
52 /// A payload that cannot be written ends the track: a log missing a record is not the lossless
53 /// log this mode promises, so the failure is surfaced rather than papered over with a second
54 /// group. The group is aborted rather than closed cleanly, so a consumer sees the failure
55 /// instead of a log that merely looks complete. Every later append fails on the closed track.
56 ///
57 /// Returns the frame's encoded size.
58 pub fn append(&mut self, payload: impl Into<Timed<Bytes>>) -> Result<usize> {
59 self.inner.lock().unwrap().append(payload.into())
60 }
61
62 /// Finish the track, closing the group.
63 pub fn finish(&mut self) -> Result<()> {
64 self.inner.lock().unwrap().finish()
65 }
66}
67
68/// Shared publishing state behind [`Producer`]'s `Arc<Mutex>`.
69struct Inner {
70 track: moq_net::track::Producer,
71
72 /// Opened on the first append and never rolled.
73 group: Option<moq_net::group::Producer>,
74
75 /// The DEFLATE encoder, one window for the whole group, `Some` while compressing.
76 flate: Option<moq_flate::Encoder>,
77}
78
79impl Inner {
80 fn append(&mut self, payload: Timed<Bytes>) -> Result<usize> {
81 let timestamp = payload.at.unwrap_or_else(moq_net::Timestamp::now);
82 let payload = payload.value;
83
84 // A payload no consumer could decode is as terminal as one the track rejects: the log is
85 // missing a record either way, and carrying on would present that gap as a complete log.
86 // Checked before the group is opened, so nothing is published, and routed through the same
87 // abort so a reader sees the failure rather than a clean end.
88 if self.flate.is_some() && payload.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE {
89 self.abort(moq_net::Error::FrameTooLarge);
90 return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into());
91 }
92
93 // Open the group before compressing: a failure here must not leave the window ahead of a
94 // consumer that never received the frame.
95 if self.group.is_none() {
96 self.group = Some(self.track.append_group()?);
97 }
98
99 let payload = match self.flate.as_mut() {
100 Some(flate) => flate.frame(&payload),
101 None => payload,
102 };
103
104 let size = payload.len();
105 let group = self.group.as_mut().expect("a group is open");
106 let Err(err) = group.write_frame(timestamp, payload) else {
107 return Ok(size);
108 };
109
110 // The payload never reached the wire, so the log has a hole in it, which is not the lossless
111 // log this mode promises. Continuing into a second group would hand consumers a gap dressed up
112 // as a complete log, so end the track and let the caller start a new one. This is also what
113 // keeps "a stream is one group" a real invariant rather than the usual case.
114 //
115 // Abort the track rather than finishing it: a clean close drains a consumer to `None`, which
116 // is exactly what a completed log looks like, so a truncated log would be indistinguishable
117 // from a whole one. Aborting the *track* is what a subscriber observes; aborting only the
118 // group drops it from the cache and the consumer still reads a clean end.
119 self.abort(err.clone());
120
121 Err(err.into())
122 }
123
124 /// End the track with an error, so a consumer sees the failure rather than a clean end.
125 fn abort(&mut self, err: moq_net::Error) {
126 // Abort the group with the same error first. `track::Producer::abort` deliberately leaves an
127 // already-pulled `group::Consumer` independent, so dropping our handle would hand a reader
128 // sitting in the group a generic `Dropped` instead of the failure that ended the log.
129 if let Some(group) = self.group.take() {
130 let _ = group.abort(err.clone());
131 }
132
133 // Abort through a clone, since aborting consumes a handle and the state is shared. Keeping
134 // ours means `consume` still hands back a subscriber, which is how a reader learns the log
135 // ended badly rather than cleanly.
136 let _ = self.track.clone().abort(err);
137 }
138
139 fn finish(&mut self) -> Result<()> {
140 // Finalize both independently rather than short-circuiting on the group. Returning early
141 // would leave the track open with `group` already taken, so a later append would open a
142 // second group, and (with compression) write into it from a window the consumer never
143 // received. That is exactly the split log ending the track exists to prevent.
144 let group = match self.group.take() {
145 Some(group) => group.finish(),
146 None => Ok(()),
147 };
148 let track = self.track.finish();
149
150 group?;
151 track?;
152 Ok(())
153 }
154}