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}