Skip to main content

moq_binary/snapshot/
producer.rs

1//! Publishing a binary value 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 a binary value over a track, one value per group.
13///
14/// Each [`update`](Self::update) rolls a new group holding the whole value, so a consumer only ever
15/// needs the newest group and older ones are dropped. For a log where every payload survives, use
16/// [`stream`](crate::stream) instead.
17///
18/// Cheaply clonable: clones share one underlying track, like other MoQ producers.
19#[derive(Clone)]
20pub struct Producer {
21	inner: Arc<Mutex<Inner>>,
22}
23
24impl Producer {
25	/// Create a producer that publishes to the given track.
26	pub fn new(track: moq_net::track::Producer, config: Config) -> Self {
27		Self {
28			inner: Arc::new(Mutex::new(Inner {
29				track,
30				compression: config.compression.is_deflate(),
31			})),
32		}
33	}
34
35	/// Create a subscriber for the underlying track.
36	pub fn consume(&self) -> moq_net::track::Subscriber {
37		self.inner.lock().unwrap().track.subscribe(None)
38	}
39
40	/// Whether any consumer for the underlying track currently exists.
41	///
42	/// The demand signal for a producer serving on request: an unused track is cached state nobody is
43	/// watching, safe to drop and recreate on the next request.
44	pub fn is_used(&self) -> bool {
45		self.inner.lock().unwrap().track.is_used()
46	}
47
48	/// Publish a new value, superseding the previous one, and return the frame's encoded size.
49	///
50	/// Unlike [`moq-json`](https://docs.rs/moq-json), an identical value is republished rather than
51	/// skipped: comparing two opaque blobs costs a full scan, and only the caller knows whether its
52	/// bytes changed.
53	pub fn update(&mut self, payload: impl Into<Timed<Bytes>>) -> Result<usize> {
54		self.inner.lock().unwrap().update(payload.into())
55	}
56
57	/// Finish the track.
58	pub fn finish(&mut self) -> Result<()> {
59		self.inner.lock().unwrap().finish()
60	}
61}
62
63/// Shared publishing state behind [`Producer`]'s `Arc<Mutex>`.
64struct Inner {
65	track: moq_net::track::Producer,
66	compression: bool,
67}
68
69impl Inner {
70	fn update(&mut self, payload: Timed<Bytes>) -> Result<usize> {
71		let timestamp = payload.at.unwrap_or_else(moq_net::Timestamp::now);
72		let payload = payload.value;
73
74		// One frame per group, so the window spans a single value and starts cold every time.
75		let payload = match self.compression {
76			true => {
77				// Compression can take a large value under the group's frame limit, but every consumer
78				// decodes with moq-flate's default output cap, so publishing past it would advertise a
79				// value that always fails to read. Reject it here instead.
80				if payload.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE {
81					return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into());
82				}
83				moq_flate::Encoder::new().frame(&payload)
84			}
85			false => payload,
86		};
87
88		// Check before opening a group. `append_group` publishes immediately, so discovering the limit
89		// inside `write_frame` would leave an empty newest group behind: a snapshot consumer jumps to
90		// the newest, so the previous value would be lost even though this update reported an error.
91		if payload.len() as u64 > moq_net::group::MAX_CACHE_BYTES {
92			return Err(moq_net::Error::FrameTooLarge.into());
93		}
94
95		let size = payload.len();
96		let mut group = self.track.append_group()?;
97		if let Err(err) = group.write_frame(timestamp, payload) {
98			// `append_group` already published this group, and a rejected frame (too large) doesn't
99			// close the track. Dropping the handle does NOT close the group, so leaving it would strand
100			// any subscriber that advanced into it with nothing to read and no end.
101			let _ = group.finish();
102			return Err(err.into());
103		}
104
105		group.finish()?;
106		Ok(size)
107	}
108
109	fn finish(&mut self) -> Result<()> {
110		self.track.finish()?;
111		Ok(())
112	}
113}