Skip to main content

moq_json/snapshot/
encoder.rs

1//! The track-free half of snapshot publishing: values in, frame payloads out.
2
3use std::cell::RefCell;
4use std::marker::PhantomData;
5use std::sync::OnceLock;
6
7use bytes::Bytes;
8use serde::Serialize;
9use serde_json::Value;
10
11use crate::{Compression, Result};
12
13/// Maximum frames (snapshot + deltas) in a single group before a new snapshot is forced.
14///
15/// Kept well below moq-net's per-group frame cap so a late joiner can always read the snapshot
16/// at frame 0 before the group is evicted.
17pub(super) const MAX_DELTA_FRAMES: usize = 256;
18
19/// What an [`Encoder`] keeps of the value it last emitted.
20///
21/// A delta is a diff against the previous value, so one has to be parsed to diff against
22/// whenever deltas are possible. With `delta_ratio = 0` none ever are, and the only question
23/// an update asks of the baseline is whether the value changed at all, which the encoded
24/// bytes answer directly. The parse is deferred in that case, and a value that is only ever
25/// published never pays for one.
26enum Baseline {
27	/// Deltas are possible, so the baseline is kept parsed and ready to diff against.
28	Parsed(Value),
29
30	/// Deltas are disabled. The emitted bytes (shared with the frame payload when not
31	/// compressing) stand in for the value, parsed only if a caller reads it back.
32	Encoded {
33		bytes: Bytes,
34		parsed: OnceLock<Option<Value>>,
35	},
36}
37
38impl Baseline {
39	/// The baseline as a parsed value, parsing the encoded bytes on first use.
40	fn value(&self) -> Option<&Value> {
41		match self {
42			Self::Parsed(value) => Some(value),
43			// Serialized by us, so this parses unless the caller's `Serialize` emitted
44			// something `serde_json` will not read back.
45			Self::Encoded { bytes, parsed } => parsed.get_or_init(|| serde_json::from_slice(bytes).ok()).as_ref(),
46		}
47	}
48}
49
50/// Codec options for an [`Encoder`], and so for the [`Producer`](super::Producer) wrapping one.
51///
52/// Build from [`Default`] and override fields (the struct is `#[non_exhaustive]`, so new
53/// options stay additive), or chain [`with_delta_ratio`](Self::with_delta_ratio).
54#[derive(Debug, Clone)]
55#[non_exhaustive]
56pub struct Config {
57	/// Controls how aggressively the encoder emits deltas (merge patches) instead of full snapshots.
58	///
59	/// A ratio of `0` disables deltas: every change is encoded as a new snapshot.
60	///
61	/// A positive ratio enables deltas. A new snapshot is emitted once the deltas *already written*
62	/// to the current group (excluding the snapshot frame) exceed `ratio` times the snapshot size.
63	/// The pending delta is excluded from that check, so the one that first crosses the budget
64	/// still lands before the group rolls. So `1` allows roughly one snapshot's worth of deltas before
65	/// rolling, and a larger ratio tolerates more.
66	///
67	/// When [`compression`](Self::compression) is [`Compression::Deflate`], both sides of the
68	/// comparison are measured on the *compressed* frame sizes (the real wire cost).
69	///
70	/// Defaults to `8`.
71	pub delta_ratio: u32,
72
73	/// Compress each group as one sync-flushed DEFLATE stream, so deltas reuse the snapshot as
74	/// context and shrink sharply.
75	///
76	/// [`Compression::None`] (the default) emits plaintext JSON frames, identical on the wire to an
77	/// uncompressed track. A [`Decoder`](super::Decoder) reading them must set the same
78	/// [`compression`](Self::compression).
79	pub compression: Compression,
80}
81
82impl Config {
83	/// Set [`delta_ratio`](Self::delta_ratio) (a builder, since the struct is `#[non_exhaustive]`).
84	pub fn with_delta_ratio(mut self, delta_ratio: u32) -> Self {
85		self.delta_ratio = delta_ratio;
86		self
87	}
88}
89
90impl Default for Config {
91	fn default() -> Self {
92		Self {
93			delta_ratio: 8,
94			compression: Compression::None,
95		}
96	}
97}
98
99/// One encoded frame, and the group boundary it implies.
100#[derive(Clone, Debug)]
101pub struct Encoded {
102	/// The frame payload, DEFLATE-compressed when [`Config::compression`] is [`Compression::Deflate`].
103	pub payload: Bytes,
104
105	/// Whether this frame is a full snapshot, which must open a new group.
106	///
107	/// `true` means the caller writes it as the first frame of a fresh group; `false` means it is a
108	/// merge patch that must be appended to the group the last snapshot opened. Mapping straight onto
109	/// [`moq_mux::container::Frame::keyframe`] is the point of the name.
110	///
111	/// The encoder decides this, never the caller: a value that sets a field to JSON null, or whose
112	/// root isn't an object, cannot be expressed as a merge patch at all, and the delta budget and
113	/// frame cap force a snapshot independently of what the caller wanted.
114	///
115	/// [`moq_mux::container::Frame::keyframe`]: https://docs.rs/moq-mux/latest/moq_mux/container/struct.Frame.html
116	pub keyframe: bool,
117}
118
119/// An encoded frame the caller has not yet acknowledged writing.
120///
121/// Returned by [`Encoder::update`]. Read [`payload`](Encoded::payload) and
122/// [`keyframe`](Encoded::keyframe) through the [`Deref`](std::ops::Deref) to [`Encoded`], write the
123/// frame, then [`commit`](Self::commit).
124///
125/// Dropping it uncommitted [`Encoder::reset`]s, so a frame that never reached the wire leaves the
126/// encoder resynchronizing with a fresh snapshot rather than emitting deltas against a baseline no
127/// consumer received. Note that this is a recovery, not a rollback: producing a delta payload
128/// advances the group's DEFLATE window, and that can't be undone, so a snapshot is the only sound
129/// way back. Forgetting to commit a frame that *was* written is therefore merely wasteful (one
130/// redundant snapshot), never incorrect.
131#[must_use = "the frame must be written and committed, or dropped to resynchronize the encoder"]
132pub struct Pending<'a, T> {
133	encoder: &'a mut Encoder<T>,
134	encoded: Encoded,
135	committed: bool,
136}
137
138impl<T> Pending<'_, T> {
139	/// Acknowledge that the frame reached the wire, keeping the encoder's state.
140	///
141	/// Only call this once the write has actually succeeded. Committing a frame that failed to write
142	/// is the one thing that corrupts the stream.
143	pub fn commit(mut self) {
144		self.committed = true;
145	}
146}
147
148impl<T> std::ops::Deref for Pending<'_, T> {
149	type Target = Encoded;
150
151	fn deref(&self) -> &Encoded {
152		&self.encoded
153	}
154}
155
156impl<T> Drop for Pending<'_, T> {
157	fn drop(&mut self) {
158		if !self.committed {
159			self.encoder.reset();
160		}
161	}
162}
163
164/// Encodes a JSON value into frame payloads, choosing snapshots and deltas automatically.
165///
166/// The track-free core of [`Producer`](super::Producer): it decides *what bytes go in a frame* and
167/// *where the group boundaries fall*, and leaves writing them to the caller. Reach for it when
168/// something else already owns the track, for example a
169/// [`moq_mux::container::Producer`](https://docs.rs/moq-mux/latest/moq_mux/container/struct.Producer.html)
170/// that is also managing a timeline and a catalog estimate:
171///
172/// ```ignore
173/// if let Some(frame) = encoder.update(&value)? {
174///     container.write(moq_mux::container::Frame {
175///         timestamp,
176///         duration: None,
177///         payload: frame.payload.clone(),
178///         keyframe: frame.keyframe,
179///     })?; // an early return here drops `frame`, resetting the encoder
180///     frame.commit();
181/// }
182/// ```
183///
184/// Frames must reach the wire in the order they were encoded, and a frame with
185/// [`keyframe`](Encoded::keyframe) set must open a new group: both the merge patches and the
186/// group-scoped DEFLATE window depend on it. [`update`](Self::update) hands back a [`Pending`]
187/// rather than a bare [`Encoded`] so a frame that never reaches the wire can't silently desync the
188/// encoder: dropping it uncommitted [`reset`](Self::reset)s, and the next value is encoded as a
189/// fresh snapshot. Committing a frame you failed to write is the one way to corrupt the stream.
190///
191/// If the caller cuts a group for its own reasons (a `cut`, `seek`, or discontinuity), call
192/// [`reset`](Self::reset) directly so the next value opens the new group with a snapshot.
193pub struct Encoder<T> {
194	config: Config,
195
196	/// The last encoded value, the baseline every delta is diffed against. `None` until the first
197	/// snapshot, which is what makes that first [`update`](Self::update) a keyframe.
198	last: Option<Baseline>,
199
200	/// Reused key buffers for comparing unchanged fields without per-update allocations, and the
201	/// memoized root entries that let an unchanged entry skip the baseline walk.
202	scratch: RefCell<crate::diff::Scratch>,
203
204	/// The current group's DEFLATE encoder (one window per group), `Some` while compressing.
205	flate: Option<moq_flate::Encoder>,
206
207	/// Bytes of deltas emitted into the current group, excluding the snapshot frame. Compressed
208	/// slice sizes when compressing, raw patch sizes otherwise.
209	delta_bytes: u64,
210
211	/// Reference size the delta budget is measured against: the current group's snapshot frame.
212	/// Its compressed slice size when compressing, raw otherwise.
213	snapshot_len: u64,
214
215	/// Frames emitted into the current group, snapshot included.
216	group_frames: usize,
217
218	/// The most bytes a group may hold, the cap a delta is admitted under. A field rather than the
219	/// constant so a test can shrink it to a size it can reach.
220	max_group_bytes: u64,
221
222	/// Whether the next frame has to be a full snapshot, because a frame was lost or the caller cut
223	/// the group. Kept separate from [`last`](Self::last) so a resync doesn't erase the value: that
224	/// field is also what [`Producer::modify`](super::Producer::modify) seeds an edit from, and dropping
225	/// it there would publish a document with every other field missing.
226	resync: bool,
227
228	_marker: PhantomData<fn(T)>,
229}
230
231impl<T> Encoder<T> {
232	/// Create an encoder with a cold baseline, so the first [`update`](Self::update) is a snapshot.
233	pub fn new(config: Config) -> Self {
234		Self {
235			config,
236			last: None,
237			scratch: RefCell::new(crate::diff::Scratch::memoized()),
238			flate: None,
239			delta_bytes: 0,
240			snapshot_len: 0,
241			group_frames: 0,
242			max_group_bytes: moq_net::group::MAX_CACHE_BYTES,
243			resync: false,
244			_marker: PhantomData,
245		}
246	}
247
248	/// The last encoded value, or `None` before the first snapshot.
249	///
250	/// This is the baseline the next delta is diffed against, which is what a caller editing the
251	/// value in place needs to start from.
252	///
253	/// With deltas disabled the baseline is held as the encoded bytes, so the first call parses
254	/// them; the result is cached, and callers that never read the value never pay for it.
255	pub fn value(&self) -> Option<&Value> {
256		self.last.as_ref()?.value()
257	}
258
259	/// Force the next [`update`](Self::update) to emit a full snapshot, even for an unchanged value.
260	///
261	/// Call this whenever the caller closes the current group behind the encoder's back (a
262	/// `cut`, a `seek`, a discontinuity). Without it the next value may be encoded as a delta
263	/// against a DEFLATE window and a baseline that the new group doesn't carry.
264	///
265	/// [`value`](Self::value) survives: the snapshot republishes it in full anyway, and it is what a
266	/// caller editing in place starts from.
267	pub fn reset(&mut self) {
268		self.flate = None;
269		self.delta_bytes = 0;
270		self.snapshot_len = 0;
271		self.group_frames = 0;
272		self.resync = true;
273	}
274}
275
276impl<T: Serialize> Encoder<T> {
277	/// Encode a new value, as a snapshot or a delta.
278	///
279	/// Returns `None` when the value is unchanged from the last one encoded, so nothing needs to be
280	/// written. Otherwise the frame comes back as a [`Pending`] the caller writes and then
281	/// [`commit`](Pending::commit)s; dropping it uncommitted resynchronizes the encoder.
282	pub fn update(&mut self, value: &T) -> Result<Option<Pending<'_, T>>> {
283		Ok(self.encode(value)?.map(|encoded| Pending {
284			encoder: self,
285			encoded,
286			committed: false,
287		}))
288	}
289
290	/// Encode a new value into a bare frame, advancing the encoder's state.
291	///
292	/// The state change is what [`Pending`] guards, so this stays private: every caller goes through
293	/// [`update`](Self::update) and has to say whether the frame reached the wire.
294	fn encode(&mut self, value: &T) -> Result<Option<Encoded>> {
295		// A lost frame, or a group the caller cut, leaves the consumer's state unknown. Re-seed with a
296		// full snapshot even when the value is unchanged, since the frame that carried it may never
297		// have landed.
298		if self.resync {
299			return self.snapshot(value).map(Some);
300		}
301
302		// With deltas disabled there is nothing to diff, so the only question is whether the value
303		// changed: compare the encodings rather than parsing a baseline to diff against. The bytes
304		// are handed straight to the snapshot when it did change, so an update still serializes
305		// `T` exactly once.
306		if let Some(Baseline::Encoded { bytes, .. }) = self.last.as_ref() {
307			let bytes = bytes.clone();
308			let next = serde_json::to_vec(value)?;
309			if next.as_slice() == bytes.as_ref() {
310				return Ok(None);
311			}
312			return self.snapshot_encoded(next).map(Some);
313		}
314
315		// The first update has no baseline to diff against, so it seeds the stream with a snapshot.
316		let Some(Baseline::Parsed(last)) = self.last.as_ref() else {
317			return self.snapshot(value).map(Some);
318		};
319
320		// Diff straight off `T`, without building a full `Value` for the new value first.
321		let crate::diff::PatchBytes { patch, forced_snapshot } =
322			crate::diff::bytes(last, value, &self.scratch).map_err(crate::Error::Json)?;
323
324		// An empty object patch with no forced null means the value is unchanged: encode nothing.
325		if !forced_snapshot && patch.is_empty() {
326			self.scratch.get_mut().commit_memo();
327			return Ok(None);
328		}
329
330		// A forced snapshot (a genuine null, or a non-object root) or an exhausted delta budget starts a
331		// new group; otherwise the change rides as a delta in the open one.
332		if forced_snapshot || !self.delta_allowed() {
333			return self.snapshot(value).map(Some);
334		}
335
336		// Compress into the per-group window only now, for a frame we are committed to emitting.
337		let bytes = Bytes::from(patch);
338
339		// Same cap as a snapshot, on the patch's plaintext: a delta that decompresses past the
340		// consumer's limit makes the whole group unreadable, since there is no keyframe after it to
341		// resynchronize on. Rejecting here leaves the encoder to reset and the group as it was.
342		if self.config.compression.is_deflate() && bytes.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE {
343			return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into());
344		}
345		let payload = match self.flate.as_mut() {
346			Some(flate) => flate.frame(&bytes),
347			None => bytes.clone(),
348		};
349
350		// A delta is only readable while the group still holds the snapshot it applies to.
351		// Admitting a patch that pushes the group past that budget would abort it
352		// (`GroupTooLarge`), leaving a late subscriber with no value. Roll a fresh snapshot
353		// instead, which is cheap next to losing the value.
354		//
355		// Measured on the encoded payload rather than the plaintext: a sync-flushed DEFLATE frame can
356		// come out slightly larger than its input, so the plaintext is not an upper bound. Compressing
357		// first advances the window, but [`Self::snapshot`] opens a fresh one, so an over-budget delta
358		// costs only the wasted compression.
359		if self.snapshot_len + self.delta_bytes + payload.len() as u64 > self.max_group_bytes {
360			return self.snapshot(value).map(Some);
361		}
362
363		self.delta_bytes += payload.len() as u64;
364		self.group_frames += 1;
365
366		// Fold the delta into the baseline so the next diff is against the value we just encoded.
367		// Reaching a delta means `delta_allowed`, which means a non-zero ratio, which is what keeps
368		// the baseline parsed.
369		let Some(Baseline::Parsed(last)) = self.last.as_mut() else {
370			unreachable!("a parsed snapshot precedes any delta")
371		};
372		crate::merge::apply_generated_bytes(last, &bytes)?;
373		self.scratch.get_mut().commit_memo();
374
375		Ok(Some(Encoded {
376			payload,
377			keyframe: false,
378		}))
379	}
380
381	/// Whether the current change may ride as a delta in the open group.
382	///
383	/// The budget gate measures the deltas *already emitted* (excluding the frame about to land)
384	/// against the group's snapshot frame. Both are compressed sizes when compressing and raw
385	/// otherwise, so the comparison is like-for-like. Because the pending frame is excluded, the delta
386	/// that tips the group past `ratio * snapshot` still lands: a group overshoots by at most one delta
387	/// before rolling.
388	fn delta_allowed(&self) -> bool {
389		let ratio = u64::from(self.config.delta_ratio);
390		ratio != 0
391			&& self.group_frames > 0
392			&& self.group_frames < MAX_DELTA_FRAMES
393			&& self.delta_bytes <= ratio * self.snapshot_len
394	}
395
396	/// Encode a full snapshot of `value`, opening a new group and reseeding the baseline.
397	fn snapshot(&mut self, value: &T) -> Result<Encoded> {
398		// Serialize directly from `value` so the snapshot frame preserves the type's own field order,
399		// keeping the wire bytes identical to serializing `T` straight to a frame.
400		let snapshot = serde_json::to_vec(value)?;
401		self.snapshot_encoded(snapshot)
402	}
403
404	/// [`snapshot`](Self::snapshot) for a value that is already serialized, so an update that
405	/// encoded `T` to compare it against a byte baseline does not encode it a second time.
406	fn snapshot_encoded(&mut self, snapshot: Vec<u8>) -> Result<Encoded> {
407		// Every consumer decodes with moq-flate's default output cap, so a value past it would be
408		// unreadable however small it compresses to. Reject it before anything is published, so the
409		// previously published value stands rather than being superseded by one nothing can read.
410		if self.config.compression.is_deflate() && snapshot.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE {
411			return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into());
412		}
413
414		// With deltas possible, read the baseline back out of those same bytes rather than
415		// serializing `value` a second time, so the baseline IS the emitted snapshot by
416		// construction. A `Serialize` impl reading a clock or interior mutable state would otherwise
417		// seed the baseline with a value no consumer ever received, and every later delta would
418		// rebase them onto it. `T` is also only visited once, which is what a caller with an
419		// expensive or effectful `Serialize` pays for.
420		//
421		// That trades a second walk of `T` for a parse of the bytes, so it is not automatically
422		// cheaper than `to_value` (see the `baseline` benchmark); consistency is the reason. With
423		// deltas off there is no diff to rebase and no reason to pay it at all.
424		//
425		// Every fallible step runs before any state changes, so a failure leaves the encoder exactly
426		// as it was rather than half-advanced with no frame to show for it.
427		let snapshot = Bytes::from(snapshot);
428		let last = if self.config.delta_ratio == 0 {
429			// No delta will ever diff against this, so hold the bytes instead. Uncompressed, they are
430			// the same allocation the payload carries, so the baseline costs a refcount.
431			Baseline::Encoded {
432				bytes: snapshot.clone(),
433				parsed: OnceLock::new(),
434			}
435		} else {
436			Baseline::Parsed(serde_json::from_slice(&snapshot)?)
437		};
438
439		// Open a fresh per-group encoder (cold window) and compress the snapshot as frame 0, recording
440		// its wire size as the delta anchor.
441		let (payload, flate) = match self.config.compression {
442			Compression::Deflate => {
443				let mut flate = moq_flate::Encoder::new();
444				let payload = flate.frame(&snapshot);
445				(payload, Some(flate))
446			}
447			Compression::None => (snapshot, None),
448		};
449
450		self.snapshot_len = payload.len() as u64;
451		self.delta_bytes = 0;
452		self.group_frames = 1;
453		self.flate = flate;
454		self.last = Some(last);
455		self.resync = false;
456		// Seeded from the snapshot rather than a diff, so no root entry is memoized against it yet.
457		self.scratch.get_mut().clear_memo();
458
459		Ok(Encoded {
460			payload,
461			keyframe: true,
462		})
463	}
464}
465
466#[cfg(test)]
467mod test {
468	use super::*;
469	use serde_json::json;
470
471	#[test]
472	fn duplicate_serialized_keys_are_refused() {
473		use serde::ser::SerializeMap;
474		struct Duplicate {
475			duplicate: bool,
476		}
477		impl Serialize for Duplicate {
478			fn serialize<S: serde::Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
479				let mut map = serializer.serialize_map(Some(2 + usize::from(self.duplicate)))?;
480				map.serialize_entry("a", &1)?;
481				map.serialize_entry("b", &2)?;
482				if self.duplicate {
483					map.serialize_entry("a", &3)?;
484				}
485				map.end()
486			}
487		}
488		let mut encoder = Encoder::<Duplicate>::new(Config::default());
489		encoder
490			.update(&Duplicate { duplicate: false })
491			.unwrap()
492			.unwrap()
493			.commit();
494		let err = encoder.encode(&Duplicate { duplicate: true }).unwrap_err();
495		assert!(err.to_string().contains("duplicate JSON object key"));
496	}
497
498	/// Encode a sequence of values, committing each frame, and return `(keyframe, payload_len)` per
499	/// emitted frame.
500	fn encode(config: Config, values: &[Value]) -> Vec<(bool, usize)> {
501		let mut encoder = Encoder::<Value>::new(config);
502		let mut out = Vec::new();
503		for value in values {
504			if let Some(frame) = encoder.update(value).unwrap() {
505				out.push((frame.keyframe, frame.payload.len()));
506				frame.commit();
507			}
508		}
509		out
510	}
511
512	/// Encode one value and commit it, returning the frame.
513	fn commit(encoder: &mut Encoder<Value>, value: &Value) -> Option<Encoded> {
514		let frame = encoder.update(value).unwrap()?;
515		let encoded = Encoded {
516			payload: frame.payload.clone(),
517			keyframe: frame.keyframe,
518		};
519		frame.commit();
520		Some(encoded)
521	}
522
523	#[test]
524	fn first_update_is_a_keyframe() {
525		let frames = encode(Config::default(), &[json!({ "a": 1 })]);
526		assert_eq!(frames.len(), 1);
527		assert!(frames[0].0);
528	}
529
530	#[test]
531	fn unchanged_value_encodes_nothing() {
532		let frames = encode(Config::default(), &[json!({ "a": 1 }), json!({ "a": 1 })]);
533		assert_eq!(frames.len(), 1);
534	}
535
536	#[test]
537	fn changes_ride_as_deltas() {
538		let frames = encode(
539			Config::default().with_delta_ratio(100),
540			&[
541				json!({ "a": 1, "b": 1 }),
542				json!({ "a": 1, "b": 2 }),
543				json!({ "a": 1, "b": 3 }),
544			],
545		);
546		assert_eq!(frames.iter().map(|f| f.0).collect::<Vec<_>>(), vec![true, false, false]);
547	}
548
549	#[test]
550	fn deltas_off_forces_a_keyframe_per_change() {
551		let frames = encode(
552			Config::default().with_delta_ratio(0),
553			&[json!({ "a": 1 }), json!({ "a": 2 })],
554		);
555		assert_eq!(frames.iter().map(|f| f.0).collect::<Vec<_>>(), vec![true, true]);
556	}
557
558	/// Deltas off keeps the baseline as bytes rather than a parsed value, so the unchanged check
559	/// runs on the encoding. It still has to suppress a republish, or every stats tick would
560	/// re-emit an identical frame.
561	#[test]
562	fn deltas_off_still_skips_an_unchanged_value() {
563		let frames = encode(
564			Config::default().with_delta_ratio(0),
565			&[json!({ "a": 1 }), json!({ "a": 1 }), json!({ "a": 1 })],
566		);
567		assert_eq!(frames.len(), 1);
568	}
569
570	/// Field order is part of the encoding, so a byte baseline only answers "unchanged" correctly
571	/// because `T` serializes deterministically. Same keys, different values, must still emit.
572	#[test]
573	fn deltas_off_detects_a_change_under_the_same_keys() {
574		let frames = encode(
575			Config::default().with_delta_ratio(0),
576			&[json!({ "a": 1, "b": 2 }), json!({ "a": 1, "b": 3 })],
577		);
578		assert_eq!(frames.len(), 2);
579	}
580
581	/// The byte baseline is parsed on demand, so `value` (and so `Producer::modify`, which seeds an
582	/// edit from it) keeps working with deltas off. Dropping the baseline instead would make
583	/// `modify` start from `T::default()` and publish a document with every other field missing.
584	#[test]
585	fn deltas_off_still_exposes_the_value() {
586		let mut encoder = Encoder::<Value>::new(Config::default().with_delta_ratio(0));
587		assert_eq!(encoder.value(), None);
588
589		commit(&mut encoder, &json!({ "a": 1, "b": 2 })).unwrap();
590		assert_eq!(encoder.value(), Some(&json!({ "a": 1, "b": 2 })));
591
592		commit(&mut encoder, &json!({ "a": 1, "b": 3 })).unwrap();
593		assert_eq!(encoder.value(), Some(&json!({ "a": 1, "b": 3 })));
594	}
595
596	/// Compressing shares no allocation between the baseline and the payload, so the byte baseline
597	/// has to hold the plaintext rather than the compressed frame.
598	#[test]
599	fn deltas_off_while_compressing_keeps_the_plaintext_baseline() {
600		let mut config = Config::default().with_delta_ratio(0);
601		config.compression = Compression::Deflate;
602
603		let mut encoder = Encoder::<Value>::new(config);
604		commit(&mut encoder, &json!({ "a": 1 })).unwrap();
605		assert_eq!(encoder.value(), Some(&json!({ "a": 1 })));
606		assert!(commit(&mut encoder, &json!({ "a": 1 })).is_none());
607	}
608
609	/// A value the caller might reasonably expect to be a delta, but that merge patch can't express:
610	/// setting a field to JSON null reads as a key deletion. The encoder has to override the caller
611	/// here, which is why `keyframe` is a return value rather than a parameter.
612	#[test]
613	fn a_null_field_forces_a_keyframe() {
614		let frames = encode(
615			Config::default().with_delta_ratio(100),
616			&[json!({ "a": 1, "b": 1 }), json!({ "a": 1, "b": null })],
617		);
618		assert_eq!(frames.iter().map(|f| f.0).collect::<Vec<_>>(), vec![true, true]);
619	}
620
621	/// Same story for a root that isn't an object: there is no recursive merge patch for it.
622	#[test]
623	fn a_non_object_root_forces_a_keyframe() {
624		let frames = encode(
625			Config::default().with_delta_ratio(100),
626			&[json!({ "a": 1 }), json!([1, 2, 3])],
627		);
628		assert_eq!(frames.iter().map(|f| f.0).collect::<Vec<_>>(), vec![true, true]);
629	}
630
631	#[test]
632	fn frame_cap_forces_a_keyframe() {
633		let values: Vec<Value> = (0..=MAX_DELTA_FRAMES).map(|n| json!({ "n": n })).collect();
634		let frames = encode(Config::default().with_delta_ratio(1_000_000), &values);
635
636		// The snapshot plus MAX_DELTA_FRAMES - 1 deltas fill the group, then the cap rolls it.
637		assert_eq!(frames.len(), MAX_DELTA_FRAMES + 1);
638		assert_eq!(frames.iter().filter(|f| f.0).count(), 2);
639		assert!(frames[MAX_DELTA_FRAMES].0);
640	}
641
642	/// A caller that cuts the group behind the encoder's back has to say so, or the next value would
643	/// be a delta against a window and a baseline the new group never carried.
644	#[test]
645	fn reset_forces_the_next_update_to_be_a_keyframe() {
646		let mut encoder = Encoder::<Value>::new(Config::default().with_delta_ratio(100));
647		assert!(commit(&mut encoder, &json!({ "a": 1 })).unwrap().keyframe);
648		assert!(!commit(&mut encoder, &json!({ "a": 2 })).unwrap().keyframe);
649
650		encoder.reset();
651		assert!(commit(&mut encoder, &json!({ "a": 3 })).unwrap().keyframe);
652	}
653
654	/// A frame the caller never wrote must not leave the encoder emitting deltas against a baseline
655	/// no consumer received. Dropping the [`Pending`] uncommitted is what a failed write looks like,
656	/// and it has to resynchronize on its own: a caller cannot be relied on to remember.
657	#[test]
658	fn an_uncommitted_frame_resynchronizes_the_encoder() {
659		let mut encoder = Encoder::<Value>::new(Config::default().with_delta_ratio(100));
660		commit(&mut encoder, &json!({ "a": 1 })).unwrap();
661
662		// The caller wrote this one and said so, so the next value can still ride as a delta.
663		commit(&mut encoder, &json!({ "a": 2 })).unwrap();
664
665		// This one fails to write, so the caller drops it without committing.
666		drop(encoder.update(&json!({ "a": 3 })).unwrap().expect("a delta"));
667
668		// The next value opens a new group with a full snapshot rather than patching a state the
669		// consumer never reached.
670		let recovered = commit(&mut encoder, &json!({ "a": 4 })).expect("a resynchronizing snapshot");
671		assert!(recovered.keyframe);
672		assert_eq!(
673			serde_json::from_slice::<Value>(&recovered.payload).unwrap(),
674			json!({ "a": 4 }),
675			"the snapshot carries the whole value, not a patch"
676		);
677	}
678
679	/// The same recovery when the very first frame is lost: the encoder must not treat the value as
680	/// already published and skip it as unchanged.
681	#[test]
682	fn an_uncommitted_first_frame_is_reencoded() {
683		let mut encoder = Encoder::<Value>::new(Config::default());
684		drop(encoder.update(&json!({ "a": 1 })).unwrap().expect("a snapshot"));
685
686		let retried = commit(&mut encoder, &json!({ "a": 1 })).expect("the same value, re-encoded");
687		assert!(retried.keyframe);
688	}
689
690	/// A reset value is republished even when it matches the last one encoded: the new group has to
691	/// open with a snapshot, so "unchanged" can't mean "write nothing" there.
692	#[test]
693	fn reset_republishes_an_unchanged_value() {
694		let mut encoder = Encoder::<Value>::new(Config::default());
695		commit(&mut encoder, &json!({ "a": 1 })).unwrap();
696
697		encoder.reset();
698		assert!(
699			commit(&mut encoder, &json!({ "a": 1 }))
700				.expect("a fresh snapshot")
701				.keyframe
702		);
703	}
704
705	/// A sync-flushed DEFLATE frame can come out larger than its input, so the plaintext is not an
706	/// upper bound on what lands in the group. A patch that fits the budget by its plaintext but not
707	/// by its encoded size would otherwise slip through the gate, overflow the group, and evict the
708	/// snapshot a late joiner needs. Matches the JS `a compressed delta is gated on its encoded size`.
709	#[test]
710	fn a_compressed_delta_is_gated_on_its_encoded_size() {
711		let value = json!({ "v": "x".repeat(1000) });
712		let patched = json!({ "v": "x".repeat(1000), "q": "a" });
713		let plaintext = serde_json::to_vec(&json!({ "q": "a" })).unwrap().len();
714		let config = Config {
715			delta_ratio: 100,
716			compression: Compression::Deflate,
717		};
718
719		// Measure the frames under the default budget, which admits the delta.
720		let mut probe = Encoder::<Value>::new(config.clone());
721		let snapshot = commit(&mut probe, &value).expect("a snapshot");
722		let delta = commit(&mut probe, &patched).expect("a delta");
723		assert!(!delta.keyframe);
724		assert!(
725			delta.payload.len() > plaintext,
726			"the encoded delta ({}) must outgrow its plaintext ({plaintext}) for this to test anything",
727			delta.payload.len()
728		);
729
730		// A budget with room for the plaintext patch but not the encoded one.
731		let mut encoder = Encoder::<Value>::new(config);
732		encoder.max_group_bytes = (snapshot.payload.len() + plaintext) as u64;
733		assert!(commit(&mut encoder, &value).unwrap().keyframe);
734
735		// The patch rolls a fresh snapshot instead of joining the first group.
736		let rolled = commit(&mut encoder, &patched).expect("a rolled snapshot");
737		assert!(rolled.keyframe);
738		let mut decoder = moq_flate::Decoder::new();
739		assert_eq!(
740			serde_json::from_slice::<Value>(&decoder.frame(&rolled.payload).unwrap()).unwrap(),
741			patched,
742			"the snapshot carries the whole value, not a patch"
743		);
744	}
745
746	#[test]
747	fn compressed_deltas_reuse_the_group_window() {
748		let phrase = "Media over QUIC delivers real-time latency at massive scale";
749		let frames = encode(
750			Config {
751				delta_ratio: 100,
752				compression: Compression::Deflate,
753			},
754			&[json!({ "note": phrase }), json!({ "note": phrase, "echo": phrase })],
755		);
756
757		// The raw patch repeats the whole phrase; compressed against the window it's a fraction.
758		let raw = serde_json::to_vec(&json!({ "echo": phrase })).unwrap().len();
759		assert_eq!(frames.len(), 2);
760		assert!(
761			frames[1].1 < raw / 2,
762			"windowed delta {} vs raw patch {raw}",
763			frames[1].1
764		);
765	}
766
767	/// A value whose serialization changes on every call, standing in for a `Serialize` impl backed by
768	/// a clock, an atomic, or interior mutable state.
769	struct Ticking(std::cell::Cell<u32>);
770
771	impl serde::Serialize for Ticking {
772		fn serialize<S: serde::Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
773			use serde::ser::SerializeMap;
774
775			let n = self.0.get();
776			self.0.set(n + 1);
777
778			let mut map = serializer.serialize_map(Some(1))?;
779			map.serialize_entry("n", &n)?;
780			map.end()
781		}
782	}
783
784	/// The snapshot frame and the baseline must come from a single pass over the value. Serializing
785	/// twice costs a second traversal, and for a value like this one it seeds the baseline with
786	/// something no consumer ever received, so every later delta rebases them onto a phantom state.
787	#[test]
788	fn a_snapshot_serializes_its_value_once() {
789		let value = Ticking(std::cell::Cell::new(0));
790		let mut encoder = Encoder::<Ticking>::new(Config::default());
791		let payload = {
792			let frame = encoder.update(&value).unwrap().expect("a snapshot");
793			let payload = frame.payload.clone();
794			frame.commit();
795			payload
796		};
797
798		assert_eq!(value.0.get(), 1, "the value should be serialized exactly once");
799
800		let emitted: Value = serde_json::from_slice(&payload).unwrap();
801		assert_eq!(emitted, json!({ "n": 0 }));
802		assert_eq!(encoder.value(), Some(&emitted), "the baseline must be what was emitted");
803	}
804
805	/// A root entry the memo has not seen yet is diffed from the bytes the memo recorded, not
806	/// serialized again: a second pass could disagree with the first, leaving the memo describing a
807	/// value the baseline never held.
808	#[test]
809	fn a_delta_serializes_each_entry_once() {
810		let value = std::collections::BTreeMap::from([("row", Ticking(std::cell::Cell::new(0)))]);
811		let mut encoder = Encoder::new(Config::default().with_delta_ratio(100));
812		encoder.update(&value).unwrap().expect("a snapshot").commit();
813
814		let frame = encoder.update(&value).unwrap().expect("a delta");
815		assert!(!frame.keyframe);
816		let emitted: Value = serde_json::from_slice(&frame.payload).unwrap();
817		frame.commit();
818
819		assert_eq!(
820			value["row"].0.get(),
821			2,
822			"each update should serialize the entry exactly once"
823		);
824		assert_eq!(emitted, json!({ "row": { "n": 1 } }));
825		assert_eq!(encoder.value(), Some(&emitted), "the baseline must be what was emitted");
826	}
827
828	/// A key repeated below the root is refused whether the memo meets it in a new entry or in a
829	/// value replaced wholesale, as the value diff refuses it. Letting one into the memo would pair
830	/// the repeats by position, where the consumer keeps the last.
831	#[test]
832	fn a_repeated_nested_key_is_refused_through_the_memo() {
833		use serde::ser::SerializeMap;
834
835		/// `{"row": {"o": ..}}`, where `o` is `1` or an object that repeats a key.
836		struct Doc {
837			repeat: bool,
838		}
839		struct Row<'a>(&'a Doc);
840		struct Repeat;
841
842		impl Serialize for Repeat {
843			fn serialize<S: serde::Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
844				let mut map = serializer.serialize_map(Some(2))?;
845				map.serialize_entry("x", &1)?;
846				map.serialize_entry("x", &2)?;
847				map.end()
848			}
849		}
850		impl Serialize for Row<'_> {
851			fn serialize<S: serde::Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
852				let mut map = serializer.serialize_map(Some(1))?;
853				match self.0.repeat {
854					true => map.serialize_entry("o", &Repeat)?,
855					false => map.serialize_entry("o", &1)?,
856				}
857				map.end()
858			}
859		}
860		impl Serialize for Doc {
861			fn serialize<S: serde::Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
862				let mut map = serializer.serialize_map(Some(1))?;
863				map.serialize_entry("row", &Row(self))?;
864				map.end()
865			}
866		}
867
868		let config = Config::default().with_delta_ratio(100);
869		let (plain, repeat) = (Doc { repeat: false }, Doc { repeat: true });
870
871		// A new entry: the first diff after a snapshot has nothing memoized yet.
872		let mut encoder = Encoder::<Doc>::new(config.clone());
873		encoder.update(&plain).unwrap().expect("a snapshot").commit();
874		let err = encoder.encode(&repeat).unwrap_err();
875		assert!(err.to_string().contains("duplicate JSON object key"), "{err}");
876
877		// A memoized entry whose scalar becomes an object.
878		let mut encoder = Encoder::<Doc>::new(config);
879		encoder.update(&plain).unwrap().expect("a snapshot").commit();
880		assert!(encoder.update(&plain).unwrap().is_none(), "unchanged, now memoized");
881		let err = encoder.encode(&repeat).unwrap_err();
882		assert!(err.to_string().contains("duplicate JSON object key"), "{err}");
883	}
884
885	/// A root object whose entries serialize in the order given, sorted or not.
886	struct Rows(Vec<(String, Value)>);
887
888	impl Serialize for Rows {
889		fn serialize<S: serde::Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
890			serializer.collect_map(self.0.iter().map(|(key, value)| (key, value)))
891		}
892	}
893
894	/// A deterministic xorshift, so a failure replays.
895	struct Rng(u64);
896
897	impl Rng {
898		fn below(&mut self, n: u64) -> u64 {
899			self.0 ^= self.0 << 13;
900			self.0 ^= self.0 >> 7;
901			self.0 ^= self.0 << 17;
902			self.0 % n
903		}
904
905		/// A row value covering what the memo has to get right: nested objects that gain and lose
906		/// keys, values that change type, nulls in and out of arrays, and strings that look like JSON.
907		fn row(&mut self) -> Value {
908			let strings = ["plain", "q\"uote", "back\\slash", "},{\"x\":1", "null", "a:b,c"];
909			let mut row = serde_json::Map::new();
910			row.insert(
911				"n".into(),
912				match self.below(30) {
913					0 => Value::Null,
914					n => json!(n % 4),
915				},
916			);
917			if self.below(4) > 0 {
918				row.insert("s".into(), json!(strings[self.below(strings.len() as u64) as usize]));
919			}
920			let mut nested = serde_json::Map::new();
921			nested.insert("a".into(), json!(self.below(3)));
922			if self.below(3) == 0 {
923				nested.insert("b".into(), json!([self.below(2), null]));
924			}
925			if self.below(40) == 0 {
926				nested.insert("c".into(), Value::Null);
927			}
928			row.insert("o".into(), Value::Object(nested));
929			row.insert(
930				"t".into(),
931				match self.below(5) {
932					0 => json!({ "k": self.below(2) }),
933					1 => json!({}),
934					2 => json!([{ "k": null }]),
935					3 => json!(1.5 + self.below(2) as f64),
936					_ => json!("t"),
937				},
938			);
939			if self.below(60) == 0 {
940				row.insert("z".into(), Value::Null);
941			}
942			Value::Object(row)
943		}
944	}
945
946	/// The memo is a shortcut past the value diff, so it must never change a frame: every payload and
947	/// keyframe has to match an encoder diffing without it, through inserts, deletions, reorders,
948	/// shape changes, forced snapshots, and group rolls.
949	#[test]
950	fn memo_matches_the_value_diff() {
951		for (seed, compression) in [
952			(1, Compression::None),
953			(2, Compression::Deflate),
954			(3, Compression::None),
955		] {
956			let mut config = Config::default().with_delta_ratio(2);
957			config.compression = compression;
958			let mut memoized = Encoder::<Rows>::new(config.clone());
959			let mut plain = Encoder::<Rows>::new(config);
960			plain.scratch = RefCell::new(crate::diff::Scratch::default());
961
962			let mut rng = Rng(0x9E37_79B9_7F4A_7C15 ^ seed);
963			let mut rows: Vec<(String, Value)> = (0..40).map(|i| (format!("row-{i:03}"), rng.row())).collect();
964			let mut emitted = 0;
965			for tick in 0..400 {
966				for row in rows.iter_mut() {
967					if rng.below(4) == 0 {
968						row.1 = rng.row();
969					}
970				}
971				if rng.below(3) == 0 {
972					let index = rng.below(rows.len() as u64) as usize;
973					rows.remove(index);
974				}
975				if rng.below(3) == 0 {
976					rows.push((format!("row-{:03}", 40 + rng.below(40)), rng.row()));
977				}
978				rows.sort_by(|a, b| a.0.cmp(&b.0));
979				rows.dedup_by(|a, b| a.0 == b.0);
980				// Now and then, a root that stops ascending.
981				if seed == 3 && rng.below(10) == 0 {
982					let (a, b) = (
983						rng.below(rows.len() as u64) as usize,
984						rng.below(rows.len() as u64) as usize,
985					);
986					rows.swap(a, b);
987				}
988
989				let value = Rows(rows.clone());
990				let want = plain.update(&value).unwrap().map(|frame| {
991					let encoded = (*frame).clone();
992					frame.commit();
993					encoded
994				});
995				let got = memoized.update(&value).unwrap().map(|frame| {
996					let encoded = (*frame).clone();
997					frame.commit();
998					encoded
999				});
1000				match (want, got) {
1001					(None, None) => {}
1002					(Some(want), Some(got)) => {
1003						assert_eq!(got.keyframe, want.keyframe, "seed {seed} tick {tick}: keyframe");
1004						assert_eq!(got.payload, want.payload, "seed {seed} tick {tick}: payload");
1005						emitted += usize::from(!got.keyframe);
1006					}
1007					(want, got) => panic!("seed {seed} tick {tick}: {want:?} vs {got:?}"),
1008				}
1009				assert_eq!(memoized.value(), plain.value(), "seed {seed} tick {tick}: baseline");
1010			}
1011			assert!(emitted > 100, "seed {seed}: only {emitted} deltas exercised the memo");
1012		}
1013	}
1014
1015	#[test]
1016	fn value_tracks_the_baseline() {
1017		let mut encoder = Encoder::<Value>::new(Config::default().with_delta_ratio(100));
1018		assert_eq!(encoder.value(), None);
1019
1020		commit(&mut encoder, &json!({ "a": 1, "b": 1 }));
1021		commit(&mut encoder, &json!({ "a": 1, "b": 2 }));
1022
1023		// The delta was folded into the baseline, so it reflects what was actually published.
1024		assert_eq!(encoder.value(), Some(&json!({ "a": 1, "b": 2 })));
1025	}
1026}