1use std::marker::PhantomData;
20use std::ops::{Deref, DerefMut};
21use std::sync::{Arc, Mutex, MutexGuard};
22use std::task::Poll;
23
24use bytes::Bytes;
25use moq_flate::{Decoder, Encoder};
26use serde::Serialize;
27use serde::de::DeserializeOwned;
28use serde_json::Value;
29
30use crate::{Diff, Error, Result, diff};
31
32const MAX_DELTA_FRAMES: usize = 256;
37#[derive(Debug, Clone)]
42#[non_exhaustive]
43pub struct ProducerConfig {
44 pub delta_ratio: u32,
59
60 pub compression: bool,
66}
67
68impl ProducerConfig {
69 pub fn with_delta_ratio(mut self, delta_ratio: u32) -> Self {
71 self.delta_ratio = delta_ratio;
72 self
73 }
74
75 pub fn with_compression(mut self, compression: bool) -> Self {
77 self.compression = compression;
78 self
79 }
80}
81
82impl Default for ProducerConfig {
83 fn default() -> Self {
84 Self {
85 delta_ratio: 8,
86 compression: false,
87 }
88 }
89}
90
91#[derive(Debug, Clone, Default)]
96#[non_exhaustive]
97pub struct ConsumerConfig {
98 pub compression: bool,
101}
102
103impl ConsumerConfig {
104 pub fn with_compression(mut self, compression: bool) -> Self {
106 self.compression = compression;
107 self
108 }
109}
110
111pub struct Producer<T> {
116 inner: Arc<Mutex<Inner>>,
117 _marker: PhantomData<fn(T)>,
118}
119
120impl<T> Clone for Producer<T> {
121 fn clone(&self) -> Self {
122 Self {
123 inner: self.inner.clone(),
124 _marker: PhantomData,
125 }
126 }
127}
128
129impl<T> Producer<T> {
130 pub fn consume(&self) -> moq_net::track::Subscriber {
132 self.inner.lock().unwrap().track.subscribe(None)
133 }
134}
135
136impl<T: Serialize> Producer<T> {
137 pub fn new(track: moq_net::track::Producer, config: ProducerConfig) -> Self {
139 Self {
140 inner: Arc::new(Mutex::new(Inner {
141 track,
142 group: None,
143 encoder: None,
144 last: None,
145 delta_bytes: 0,
146 snapshot_len: 0,
147 group_frames: 0,
148 config,
149 })),
150 _marker: PhantomData,
151 }
152 }
153
154 pub fn update(&mut self, value: &T) -> Result<()> {
158 self.inner.lock().unwrap().update(value)
159 }
160
161 pub fn lock(&mut self) -> Guard<'_, T>
175 where
176 T: Default + DeserializeOwned,
177 {
178 let inner = self.inner.lock().unwrap();
179 let value = inner
180 .last
181 .as_ref()
182 .and_then(|last| serde_json::from_value(last.clone()).ok())
183 .unwrap_or_default();
184
185 Guard {
186 inner,
187 value,
188 dirty: false,
189 }
190 }
191
192 pub fn finish(&mut self) -> Result<()> {
194 self.inner.lock().unwrap().finish()
195 }
196}
197
198pub struct Guard<'a, T: Serialize> {
206 inner: MutexGuard<'a, Inner>,
207 value: T,
208 dirty: bool,
209}
210
211impl<T: Serialize> Guard<'_, T> {
212 pub fn commit(mut self) -> Result<()> {
217 self.publish()
218 }
219
220 fn publish(&mut self) -> Result<()> {
222 if !self.dirty {
223 return Ok(());
224 }
225 self.dirty = false;
226
227 self.inner.update(&self.value)
229 }
230}
231
232impl<T: Serialize> Deref for Guard<'_, T> {
233 type Target = T;
234
235 fn deref(&self) -> &T {
236 &self.value
237 }
238}
239
240impl<T: Serialize> DerefMut for Guard<'_, T> {
241 fn deref_mut(&mut self) -> &mut T {
242 self.dirty = true;
243 &mut self.value
244 }
245}
246
247impl<T: Serialize> Drop for Guard<'_, T> {
248 fn drop(&mut self) {
249 if let Err(err) = self.publish() {
250 tracing::warn!(%err, "failed to publish JSON value on guard drop");
251 }
252 }
253}
254
255struct Inner {
257 track: moq_net::track::Producer,
258 group: Option<moq_net::group::Producer>,
259 encoder: Option<Encoder>,
261 last: Option<Value>,
262 delta_bytes: u64,
265 snapshot_len: u64,
268 group_frames: usize,
269 config: ProducerConfig,
270}
271
272impl Inner {
273 fn update<T: Serialize>(&mut self, value: &T) -> Result<()> {
274 let Some(last) = self.last.as_ref() else {
277 return self.snapshot(value);
278 };
279
280 let Diff { patch, forced_snapshot } = diff(last, value);
282
283 if !forced_snapshot && patch.as_object().is_some_and(serde_json::Map::is_empty) {
285 return Ok(());
286 }
287
288 if forced_snapshot || !self.delta_allowed() {
291 return self.snapshot(value);
292 }
293
294 let bytes = serde_json::to_vec(&patch)?;
296 let slice = match self.encoder.as_mut() {
297 Some(encoder) => encoder.frame(&bytes),
298 None => Bytes::from(bytes),
299 };
300 let len = slice.len() as u64;
301 self.group
302 .as_mut()
303 .expect("delta_allowed guarantees an open group")
304 .write_frame(moq_net::Timestamp::now(), slice)?;
305 self.delta_bytes += len;
306 self.group_frames += 1;
307
308 json_patch::merge(self.last.as_mut().expect("a snapshot precedes any delta"), &patch);
310 Ok(())
311 }
312
313 fn delta_allowed(&self) -> bool {
321 let ratio = self.config.delta_ratio as u64;
322 ratio != 0
323 && self.group.is_some()
324 && self.group_frames < MAX_DELTA_FRAMES
325 && self.delta_bytes <= ratio * self.snapshot_len
326 }
327
328 fn snapshot<T: Serialize>(&mut self, value: &T) -> Result<()> {
330 let snapshot = serde_json::to_vec(value)?;
333
334 if let Some(mut group) = self.group.take() {
336 group.finish()?;
337 }
338
339 let mut group = self.track.append_group()?;
340
341 let (slice, encoder) = if self.config.compression {
344 let mut encoder = Encoder::new();
345 let slice = encoder.frame(&snapshot);
346 (slice, Some(encoder))
347 } else {
348 (Bytes::from(snapshot), None)
349 };
350 self.snapshot_len = slice.len() as u64;
351 group.write_frame(moq_net::Timestamp::now(), slice)?;
352 self.delta_bytes = 0;
353 self.group_frames = 1;
354 self.encoder = encoder;
355
356 if self.config.delta_ratio != 0 {
357 self.group = Some(group);
359 } else {
360 self.encoder = None;
362 group.finish()?;
363 }
364
365 self.last = Some(serde_json::to_value(value)?);
367 Ok(())
368 }
369
370 fn finish(&mut self) -> Result<()> {
371 if let Some(mut group) = self.group.take() {
372 group.finish()?;
373 }
374 self.track.finish()?;
375 Ok(())
376 }
377}
378
379pub struct Consumer<T> {
381 track: moq_net::track::Subscriber,
382 group: Option<moq_net::group::Consumer>,
383 compressed: bool,
385 decoder: Option<Decoder>,
387 current: Option<Value>,
388 frames_read: usize,
389 _marker: PhantomData<fn() -> T>,
390}
391
392impl<T: DeserializeOwned> Consumer<T> {
393 pub fn new(track: moq_net::track::Subscriber, config: ConsumerConfig) -> Self {
398 Self {
399 track,
400 group: None,
401 compressed: config.compression,
402 decoder: None,
403 current: None,
404 frames_read: 0,
405 _marker: PhantomData,
406 }
407 }
408
409 pub async fn next(&mut self) -> Result<Option<T>>
411 where
412 T: Unpin,
413 {
414 kio::wait(|waiter| self.poll_next(waiter)).await
415 }
416
417 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<T>>> {
427 let track_finished = loop {
429 match self.track.poll_next_group(waiter)? {
430 Poll::Ready(Some(group)) => {
431 self.group = Some(group);
432 self.current = None;
433 self.frames_read = 0;
434 self.decoder = None;
436 }
437 Poll::Ready(None) => break true,
438 Poll::Pending => break false,
439 }
440 };
441
442 let mut advanced = false;
447 let mut group_pending = false;
448 while let Some(group) = &mut self.group {
449 match group.poll_read_frame(waiter)? {
450 Poll::Ready(Some(frame)) => {
451 self.apply(frame.payload)?;
452 advanced = true;
453 }
454 Poll::Ready(None) => {
456 self.group = None;
457 break;
458 }
459 Poll::Pending => {
461 group_pending = true;
462 break;
463 }
464 }
465 }
466
467 if advanced {
468 return Poll::Ready(Ok(Some(self.reconstruct()?)));
470 }
471
472 if group_pending {
475 return Poll::Pending;
476 }
477
478 if track_finished {
479 Poll::Ready(Ok(None))
480 } else {
481 Poll::Pending
482 }
483 }
484
485 fn decode(&mut self, slice: Bytes) -> Result<Bytes> {
490 if !self.compressed {
491 return Ok(slice);
492 }
493
494 let decoder = self.decoder.get_or_insert_with(Decoder::new);
495 Ok(decoder.frame(&slice)?)
496 }
497
498 fn apply(&mut self, frame: Bytes) -> Result<()> {
501 let frame = self.decode(frame)?;
502 if self.frames_read == 0 {
503 self.current = Some(serde_json::from_slice(&frame)?);
504 } else {
505 let patch: Value = serde_json::from_slice(&frame)?;
506 let current = self.current.as_mut().expect("a snapshot precedes any delta");
507 json_patch::merge(current, &patch);
508 }
509 self.frames_read += 1;
510 Ok(())
511 }
512
513 fn reconstruct(&self) -> Result<T> {
521 let current = self
522 .current
523 .as_ref()
524 .expect("a value is present after applying a frame");
525
526 serde_path_to_error::deserialize(current).map_err(|err| {
527 let path = err.path().to_string();
528 match path.as_str() {
529 "." => Error::Json(err.into_inner().to_string()),
531 _ => Error::Json(format!("{}: {}", path, err.into_inner())),
532 }
533 })
534 }
535}
536
537#[cfg(test)]
538mod test {
539 use super::*;
540 use serde_json::json;
541
542 fn cfg(delta_ratio: u32) -> ProducerConfig {
544 ProducerConfig {
545 delta_ratio,
546 ..Default::default()
547 }
548 }
549
550 fn cfg_deflate(delta_ratio: u32) -> ProducerConfig {
552 ProducerConfig {
553 delta_ratio,
554 compression: true,
555 }
556 }
557
558 fn deflate_consumer(track: moq_net::track::Subscriber) -> Consumer<Value> {
560 Consumer::new(track, ConsumerConfig { compression: true })
561 }
562
563 fn producer(config: ProducerConfig) -> (Producer<Value>, moq_net::track::Subscriber) {
564 let track = moq_net::broadcast::Info::new()
565 .produce()
566 .create_track("test", None)
567 .unwrap();
568 let consumer = track.subscribe(None);
569 (Producer::new(track, config), consumer)
570 }
571
572 fn drain(track: moq_net::track::Subscriber) -> Vec<Value> {
574 drain_with(Consumer::<Value>::new(track, ConsumerConfig::default()))
575 }
576
577 fn drain_with(mut consumer: Consumer<Value>) -> Vec<Value> {
579 let waiter = kio::Waiter::noop();
580 let mut out = Vec::new();
581 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
582 out.push(value);
583 }
584 out
585 }
586
587 #[test]
588 fn deltas_off_snapshot_per_group() {
589 let (mut producer, track) = producer(cfg(0));
590 producer.update(&json!({ "a": 1 })).unwrap();
591 producer.update(&json!({ "a": 2 })).unwrap();
592 producer.finish().unwrap();
593
594 assert_eq!(track.latest(), Some(1));
597 assert_eq!(drain(track), vec![json!({ "a": 2 })]);
598 }
599
600 #[test]
601 fn live_consumer_sees_each_update() {
602 let (mut producer, track) = producer(ProducerConfig::default());
603 let mut consumer = Consumer::<Value>::new(track, ConsumerConfig::default());
604 let waiter = kio::Waiter::noop();
605
606 for n in 1..=3 {
607 producer.update(&json!({ "a": n })).unwrap();
608 match consumer.poll_next(&waiter) {
609 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": n })),
610 other => panic!("expected value, got {other:?}"),
611 }
612 }
613 }
614
615 #[test]
616 fn unchanged_value_writes_nothing() {
617 let (mut producer, track) = producer(ProducerConfig::default());
618 producer.update(&json!({ "a": 1 })).unwrap();
619 producer.update(&json!({ "a": 1 })).unwrap();
620 producer.finish().unwrap();
621
622 assert_eq!(track.latest(), Some(0));
623 assert_eq!(drain(track), vec![json!({ "a": 1 })]);
624 }
625
626 #[test]
627 fn deltas_share_one_group() {
628 let config = cfg(100);
629 let (mut producer, track) = producer(config);
630 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
631 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
632 producer.update(&json!({ "a": 1, "b": 3 })).unwrap();
633 producer.finish().unwrap();
634
635 assert_eq!(track.latest(), Some(0));
637 let values = drain(track);
638 assert_eq!(values.last().unwrap(), &json!({ "a": 1, "b": 3 }));
639 }
640
641 #[test]
642 fn tight_ratio_rolls_snapshots() {
643 let config = cfg(1);
648 let (mut producer, track) = producer(config);
649 producer.update(&json!({ "a": 1 })).unwrap(); producer.update(&json!({ "a": 2 })).unwrap(); producer.update(&json!({ "a": 3 })).unwrap(); producer.update(&json!({ "a": 4 })).unwrap(); producer.finish().unwrap();
654
655 assert_eq!(track.latest(), Some(1));
656 }
657
658 #[test]
659 fn deltas_stay_within_ratio_times_snapshot() {
660 let config = cfg(8);
666 let (mut producer, track) = producer(config);
667 for n in 0..=10 {
668 producer.update(&json!({ "n": n })).unwrap();
669 }
670 producer.finish().unwrap();
671
672 assert_eq!(track.latest(), Some(1));
674 assert_eq!(drain(track).last().unwrap(), &json!({ "n": 10 }));
675 }
676
677 #[test]
678 fn array_change_is_delta() {
679 let config = cfg(100);
680 let (mut producer, track) = producer(config);
681 producer.update(&json!({ "list": [1, 2] })).unwrap();
682 producer.update(&json!({ "list": [1, 2, 3] })).unwrap();
683 producer.finish().unwrap();
684
685 assert_eq!(track.latest(), Some(0));
687 assert_eq!(drain(track).last().unwrap(), &json!({ "list": [1, 2, 3] }));
688 }
689
690 #[test]
691 fn frame_cap_rolls_snapshot() {
692 let config = cfg(1_000_000);
693 let (mut producer, track) = producer(config);
694 for i in 0..=MAX_DELTA_FRAMES {
696 producer.update(&json!({ "n": i })).unwrap();
697 }
698 producer.finish().unwrap();
699
700 assert_eq!(track.latest(), Some(1));
702 assert_eq!(drain(track).last().unwrap(), &json!({ "n": MAX_DELTA_FRAMES }));
703 }
704
705 #[test]
706 fn late_joiner_reconstructs_from_deltas() {
707 let config = cfg(100);
708 let (mut producer, track) = producer(config);
709 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
710 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
711 producer.update(&json!({ "a": 5, "b": 2 })).unwrap();
712 producer.finish().unwrap();
713
714 assert_eq!(drain(track).last().unwrap(), &json!({ "a": 5, "b": 2 }));
716 }
717
718 #[test]
719 fn lock_composes_independent_owners() {
720 #[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
722 struct Doc {
723 #[serde(skip_serializing_if = "Option::is_none")]
724 video: Option<String>,
725 #[serde(skip_serializing_if = "Option::is_none")]
726 scte35: Option<u32>,
727 }
728
729 let track = moq_net::broadcast::Info::new()
730 .produce()
731 .create_track("test", None)
732 .unwrap();
733 let consumer = track.subscribe(None);
734 let mut producer = Producer::<Doc>::new(track, ProducerConfig::default());
735
736 producer.lock().video = Some("v1".to_string());
738
739 producer.lock().scte35 = Some(42);
741
742 let _ = producer.lock();
744
745 producer.finish().unwrap();
746
747 let mut consumer = Consumer::<Doc>::new(consumer, ConsumerConfig::default());
748 let waiter = kio::Waiter::noop();
749 let mut last = None;
750 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
751 last = Some(value);
752 }
753 assert_eq!(
754 last.unwrap(),
755 Doc {
756 video: Some("v1".to_string()),
757 scte35: Some(42),
758 }
759 );
760 }
761
762 #[test]
763 fn commit_reports_a_publish_failure() {
764 #[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
765 struct Doc {
766 a: u32,
767 }
768
769 let track = moq_net::broadcast::Info::new()
770 .produce()
771 .create_track("test", None)
772 .unwrap();
773 let mut producer = Producer::<Doc>::new(track, ProducerConfig::default());
774
775 producer.finish().unwrap();
777
778 let mut guard = producer.lock();
779 guard.a = 1;
780 assert!(matches!(guard.commit(), Err(crate::Error::Net(_))));
781 }
782
783 #[test]
784 fn commit_publishes_once() {
785 #[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
786 struct Doc {
787 a: u32,
788 }
789
790 let track = moq_net::broadcast::Info::new()
791 .produce()
792 .create_track("test", None)
793 .unwrap();
794 let consumer = track.subscribe(None);
795 let mut producer = Producer::<Doc>::new(track, cfg(0));
796
797 let mut guard = producer.lock();
798 guard.a = 1;
799 guard.commit().unwrap();
800
801 producer.finish().unwrap();
803 assert_eq!(consumer.latest(), Some(0));
804 }
805
806 #[test]
807 fn newer_group_supersedes_in_progress_reconstruction() {
808 let config = cfg(1);
811 let (mut producer, track) = producer(config);
812 let observer = producer.consume();
813 let mut consumer = Consumer::<Value>::new(track, ConsumerConfig::default());
814 let waiter = kio::Waiter::noop();
815
816 producer.update(&json!({ "a": 1 })).unwrap(); match consumer.poll_next(&waiter) {
818 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": 1 })),
819 other => panic!("expected first value, got {other:?}"),
820 }
821
822 producer.update(&json!({ "a": 2 })).unwrap(); producer.update(&json!({ "a": 3 })).unwrap(); producer.update(&json!({ "a": 4 })).unwrap(); producer.finish().unwrap();
826 assert_eq!(observer.latest(), Some(1));
827
828 let mut last = None;
830 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
831 last = Some(value);
832 }
833 assert_eq!(last.unwrap(), json!({ "a": 4 }));
834 }
835
836 #[test]
837 fn open_group_pends_after_track_finish() {
838 let mut track = moq_net::broadcast::Info::new()
841 .produce()
842 .create_track("test", None)
843 .unwrap();
844 let mut group = track.append_group().unwrap();
845 let consumer_track = track.subscribe(None);
846 track.finish().unwrap();
847
848 let mut consumer = Consumer::<Value>::new(consumer_track, ConsumerConfig::default());
849 let waiter = kio::Waiter::noop();
850
851 assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
853
854 group
855 .write_frame(
856 moq_net::Timestamp::ZERO,
857 Bytes::from(serde_json::to_vec(&json!({ "a": 1 })).unwrap()),
858 )
859 .unwrap();
860 group.finish().unwrap();
861
862 match consumer.poll_next(&waiter) {
863 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": 1 })),
864 other => panic!("expected the catalog value, got {other:?}"),
865 }
866 }
867
868 #[test]
869 fn late_joiner_collapses_backlog_to_latest() {
870 let (mut producer, track) = producer(cfg(100));
873 for n in 0..=20 {
874 producer.update(&json!({ "n": n })).unwrap();
875 }
876 producer.finish().unwrap();
877
878 assert_eq!(track.latest(), Some(0));
880 let values = drain(track);
881 assert_eq!(
882 values,
883 vec![json!({ "n": 20 })],
884 "backlog should collapse to the latest value"
885 );
886 }
887
888 #[test]
889 fn compressed_late_joiner_collapses_backlog_to_latest() {
890 let (mut producer, track) = producer(cfg_deflate(100));
892 for n in 0..=20 {
893 producer.update(&json!({ "n": n })).unwrap();
894 }
895 producer.finish().unwrap();
896
897 assert_eq!(track.latest(), Some(0));
898 let values = drain_with(deflate_consumer(track));
899 assert_eq!(
900 values,
901 vec![json!({ "n": 20 })],
902 "compressed backlog should collapse to the latest"
903 );
904 }
905
906 #[test]
907 fn compressed_snapshot_per_group_roundtrips() {
908 let (mut producer, track) = producer(cfg_deflate(0));
909 producer.update(&json!({ "a": 1 })).unwrap();
910 producer.update(&json!({ "a": 2 })).unwrap();
911 producer.finish().unwrap();
912
913 assert_eq!(track.latest(), Some(1));
915 let values = drain_with(deflate_consumer(track));
916 assert_eq!(values, vec![json!({ "a": 2 })]);
917 }
918
919 #[test]
920 fn compressed_deltas_share_one_group() {
921 let (mut producer, track) = producer(cfg_deflate(100));
922 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
923 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
924 producer.update(&json!({ "a": 1, "b": 3 })).unwrap();
925 producer.finish().unwrap();
926
927 assert_eq!(track.latest(), Some(0));
929 let values = drain_with(deflate_consumer(track));
930 assert_eq!(values.last().unwrap(), &json!({ "a": 1, "b": 3 }));
931 }
932
933 #[test]
934 fn compressed_late_joiner_reconstructs_from_deltas() {
935 let (mut producer, track) = producer(cfg_deflate(100));
936 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
937 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
938 producer.update(&json!({ "a": 5, "b": 2 })).unwrap();
939 producer.finish().unwrap();
940
941 let values = drain_with(deflate_consumer(track));
943 assert_eq!(values.last().unwrap(), &json!({ "a": 5, "b": 2 }));
944 }
945
946 #[test]
947 fn compressed_deltas_roll_on_compressed_budget() {
948 let (mut producer, track) = producer(cfg_deflate(2));
954 for n in 0..=40 {
955 producer.update(&json!({ "n": n })).unwrap();
956 }
957 producer.finish().unwrap();
958
959 assert!(
960 track.latest().unwrap() > 0,
961 "a tight ratio should roll at least one compressed group"
962 );
963 assert_eq!(drain_with(deflate_consumer(track)).last().unwrap(), &json!({ "n": 40 }));
964 }
965
966 #[test]
967 fn compression_shrinks_wire_frames() {
968 let value = json!({ "renditions": ["video".repeat(50), "video".repeat(50), "video".repeat(50)] });
970
971 let plaintext_bytes = wire_frame_len(cfg(0), &value);
972 let compressed_bytes = wire_frame_len(cfg_deflate(0), &value);
973 assert!(
974 compressed_bytes < plaintext_bytes,
975 "compressed frame {compressed_bytes} should be smaller than plaintext {plaintext_bytes}"
976 );
977 }
978
979 #[test]
980 fn compressed_deltas_reuse_window() {
981 let (mut producer, mut track) = producer(cfg_deflate(100));
984 let phrase = "Media over QUIC delivers real-time latency at massive scale";
985 producer.update(&json!({ "note": phrase })).unwrap();
986 producer.update(&json!({ "note": phrase, "echo": phrase })).unwrap();
987 producer.finish().unwrap();
988
989 let waiter = kio::Waiter::noop();
991 let Poll::Ready(Ok(Some(mut group))) = track.poll_next_group(&waiter) else {
992 panic!("expected a group");
993 };
994 let mut frames = Vec::new();
995 while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
996 frames.push(frame.payload);
997 }
998 assert_eq!(frames.len(), 2, "snapshot + one delta in a single group");
999
1000 let raw_delta = serde_json::to_vec(&json!({ "echo": phrase })).unwrap();
1002 assert!(
1003 frames[1].len() < raw_delta.len() / 2,
1004 "windowed delta {} should be far below the raw patch {}",
1005 frames[1].len(),
1006 raw_delta.len()
1007 );
1008 }
1009
1010 #[test]
1011 fn rejected_field_names_its_path() {
1012 #[derive(serde::Deserialize)]
1013 #[allow(dead_code)]
1014 struct Inner {
1015 count: u8,
1016 }
1017 #[derive(serde::Deserialize)]
1018 #[allow(dead_code)]
1019 struct Outer {
1020 inner: Inner,
1021 }
1022
1023 let (mut producer, track) = producer(cfg(0));
1024 producer.update(&json!({ "inner": { "count": 300 } })).unwrap();
1025
1026 let mut consumer = Consumer::<Outer>::new(track, ConsumerConfig::default());
1027 let Poll::Ready(Err(err)) = consumer.poll_next(&kio::Waiter::noop()) else {
1028 panic!("expected a deserialize error");
1029 };
1030
1031 assert!(err.to_string().starts_with("json: inner.count: "), "{err}");
1034 }
1035
1036 #[test]
1037 fn rejected_root_omits_the_path() {
1038 let (mut producer, track) = producer(cfg(0));
1039 producer.update(&json!("not a map")).unwrap();
1040
1041 let mut consumer = Consumer::<std::collections::BTreeMap<String, u8>>::new(track, ConsumerConfig::default());
1042 let Poll::Ready(Err(err)) = consumer.poll_next(&kio::Waiter::noop()) else {
1043 panic!("expected a deserialize error");
1044 };
1045 assert_eq!(
1046 err.to_string(),
1047 "json: invalid type: string \"not a map\", expected a map"
1048 );
1049 }
1050
1051 fn wire_frame_len(config: ProducerConfig, value: &Value) -> usize {
1053 let (mut producer, mut track) = producer(config);
1054 producer.update(value).unwrap();
1055 producer.finish().unwrap();
1056
1057 let waiter = kio::Waiter::noop();
1058 let Poll::Ready(Ok(Some(mut group))) = track.poll_next_group(&waiter) else {
1059 panic!("expected a group");
1060 };
1061 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
1063 panic!("expected a frame");
1064 };
1065 frame.payload.len()
1066 }
1067}