Skip to main content

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}