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, 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> {
516 let current = self
517 .current
518 .as_ref()
519 .expect("a value is present after applying a frame");
520 Ok(serde_json::from_value(current.clone())?)
521 }
522}
523
524#[cfg(test)]
525mod test {
526 use super::*;
527 use serde_json::json;
528
529 fn cfg(delta_ratio: u32) -> ProducerConfig {
531 ProducerConfig {
532 delta_ratio,
533 ..Default::default()
534 }
535 }
536
537 fn cfg_deflate(delta_ratio: u32) -> ProducerConfig {
539 ProducerConfig {
540 delta_ratio,
541 compression: true,
542 }
543 }
544
545 fn deflate_consumer(track: moq_net::track::Subscriber) -> Consumer<Value> {
547 Consumer::new(track, ConsumerConfig { compression: true })
548 }
549
550 fn producer(config: ProducerConfig) -> (Producer<Value>, moq_net::track::Subscriber) {
551 let track = moq_net::broadcast::Info::new()
552 .produce()
553 .create_track("test", None)
554 .unwrap();
555 let consumer = track.subscribe(None);
556 (Producer::new(track, config), consumer)
557 }
558
559 fn drain(track: moq_net::track::Subscriber) -> Vec<Value> {
561 drain_with(Consumer::<Value>::new(track, ConsumerConfig::default()))
562 }
563
564 fn drain_with(mut consumer: Consumer<Value>) -> Vec<Value> {
566 let waiter = kio::Waiter::noop();
567 let mut out = Vec::new();
568 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
569 out.push(value);
570 }
571 out
572 }
573
574 #[test]
575 fn deltas_off_snapshot_per_group() {
576 let (mut producer, track) = producer(cfg(0));
577 producer.update(&json!({ "a": 1 })).unwrap();
578 producer.update(&json!({ "a": 2 })).unwrap();
579 producer.finish().unwrap();
580
581 assert_eq!(track.latest(), Some(1));
584 assert_eq!(drain(track), vec![json!({ "a": 2 })]);
585 }
586
587 #[test]
588 fn live_consumer_sees_each_update() {
589 let (mut producer, track) = producer(ProducerConfig::default());
590 let mut consumer = Consumer::<Value>::new(track, ConsumerConfig::default());
591 let waiter = kio::Waiter::noop();
592
593 for n in 1..=3 {
594 producer.update(&json!({ "a": n })).unwrap();
595 match consumer.poll_next(&waiter) {
596 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": n })),
597 other => panic!("expected value, got {other:?}"),
598 }
599 }
600 }
601
602 #[test]
603 fn unchanged_value_writes_nothing() {
604 let (mut producer, track) = producer(ProducerConfig::default());
605 producer.update(&json!({ "a": 1 })).unwrap();
606 producer.update(&json!({ "a": 1 })).unwrap();
607 producer.finish().unwrap();
608
609 assert_eq!(track.latest(), Some(0));
610 assert_eq!(drain(track), vec![json!({ "a": 1 })]);
611 }
612
613 #[test]
614 fn deltas_share_one_group() {
615 let config = cfg(100);
616 let (mut producer, track) = producer(config);
617 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
618 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
619 producer.update(&json!({ "a": 1, "b": 3 })).unwrap();
620 producer.finish().unwrap();
621
622 assert_eq!(track.latest(), Some(0));
624 let values = drain(track);
625 assert_eq!(values.last().unwrap(), &json!({ "a": 1, "b": 3 }));
626 }
627
628 #[test]
629 fn tight_ratio_rolls_snapshots() {
630 let config = cfg(1);
635 let (mut producer, track) = producer(config);
636 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();
641
642 assert_eq!(track.latest(), Some(1));
643 }
644
645 #[test]
646 fn deltas_stay_within_ratio_times_snapshot() {
647 let config = cfg(8);
653 let (mut producer, track) = producer(config);
654 for n in 0..=10 {
655 producer.update(&json!({ "n": n })).unwrap();
656 }
657 producer.finish().unwrap();
658
659 assert_eq!(track.latest(), Some(1));
661 assert_eq!(drain(track).last().unwrap(), &json!({ "n": 10 }));
662 }
663
664 #[test]
665 fn array_change_is_delta() {
666 let config = cfg(100);
667 let (mut producer, track) = producer(config);
668 producer.update(&json!({ "list": [1, 2] })).unwrap();
669 producer.update(&json!({ "list": [1, 2, 3] })).unwrap();
670 producer.finish().unwrap();
671
672 assert_eq!(track.latest(), Some(0));
674 assert_eq!(drain(track).last().unwrap(), &json!({ "list": [1, 2, 3] }));
675 }
676
677 #[test]
678 fn frame_cap_rolls_snapshot() {
679 let config = cfg(1_000_000);
680 let (mut producer, track) = producer(config);
681 for i in 0..=MAX_DELTA_FRAMES {
683 producer.update(&json!({ "n": i })).unwrap();
684 }
685 producer.finish().unwrap();
686
687 assert_eq!(track.latest(), Some(1));
689 assert_eq!(drain(track).last().unwrap(), &json!({ "n": MAX_DELTA_FRAMES }));
690 }
691
692 #[test]
693 fn late_joiner_reconstructs_from_deltas() {
694 let config = cfg(100);
695 let (mut producer, track) = producer(config);
696 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
697 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
698 producer.update(&json!({ "a": 5, "b": 2 })).unwrap();
699 producer.finish().unwrap();
700
701 assert_eq!(drain(track).last().unwrap(), &json!({ "a": 5, "b": 2 }));
703 }
704
705 #[test]
706 fn lock_composes_independent_owners() {
707 #[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
709 struct Doc {
710 #[serde(skip_serializing_if = "Option::is_none")]
711 video: Option<String>,
712 #[serde(skip_serializing_if = "Option::is_none")]
713 scte35: Option<u32>,
714 }
715
716 let track = moq_net::broadcast::Info::new()
717 .produce()
718 .create_track("test", None)
719 .unwrap();
720 let consumer = track.subscribe(None);
721 let mut producer = Producer::<Doc>::new(track, ProducerConfig::default());
722
723 producer.lock().video = Some("v1".to_string());
725
726 producer.lock().scte35 = Some(42);
728
729 let _ = producer.lock();
731
732 producer.finish().unwrap();
733
734 let mut consumer = Consumer::<Doc>::new(consumer, ConsumerConfig::default());
735 let waiter = kio::Waiter::noop();
736 let mut last = None;
737 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
738 last = Some(value);
739 }
740 assert_eq!(
741 last.unwrap(),
742 Doc {
743 video: Some("v1".to_string()),
744 scte35: Some(42),
745 }
746 );
747 }
748
749 #[test]
750 fn commit_reports_a_publish_failure() {
751 #[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
752 struct Doc {
753 a: u32,
754 }
755
756 let track = moq_net::broadcast::Info::new()
757 .produce()
758 .create_track("test", None)
759 .unwrap();
760 let mut producer = Producer::<Doc>::new(track, ProducerConfig::default());
761
762 producer.finish().unwrap();
764
765 let mut guard = producer.lock();
766 guard.a = 1;
767 assert!(matches!(guard.commit(), Err(crate::Error::Net(_))));
768 }
769
770 #[test]
771 fn commit_publishes_once() {
772 #[derive(serde::Serialize, serde::Deserialize, Default, PartialEq, Debug)]
773 struct Doc {
774 a: u32,
775 }
776
777 let track = moq_net::broadcast::Info::new()
778 .produce()
779 .create_track("test", None)
780 .unwrap();
781 let consumer = track.subscribe(None);
782 let mut producer = Producer::<Doc>::new(track, cfg(0));
783
784 let mut guard = producer.lock();
785 guard.a = 1;
786 guard.commit().unwrap();
787
788 producer.finish().unwrap();
790 assert_eq!(consumer.latest(), Some(0));
791 }
792
793 #[test]
794 fn newer_group_supersedes_in_progress_reconstruction() {
795 let config = cfg(1);
798 let (mut producer, track) = producer(config);
799 let observer = producer.consume();
800 let mut consumer = Consumer::<Value>::new(track, ConsumerConfig::default());
801 let waiter = kio::Waiter::noop();
802
803 producer.update(&json!({ "a": 1 })).unwrap(); match consumer.poll_next(&waiter) {
805 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": 1 })),
806 other => panic!("expected first value, got {other:?}"),
807 }
808
809 producer.update(&json!({ "a": 2 })).unwrap(); producer.update(&json!({ "a": 3 })).unwrap(); producer.update(&json!({ "a": 4 })).unwrap(); producer.finish().unwrap();
813 assert_eq!(observer.latest(), Some(1));
814
815 let mut last = None;
817 while let Poll::Ready(Ok(Some(value))) = consumer.poll_next(&waiter) {
818 last = Some(value);
819 }
820 assert_eq!(last.unwrap(), json!({ "a": 4 }));
821 }
822
823 #[test]
824 fn open_group_pends_after_track_finish() {
825 let mut track = moq_net::broadcast::Info::new()
828 .produce()
829 .create_track("test", None)
830 .unwrap();
831 let mut group = track.append_group().unwrap();
832 let consumer_track = track.subscribe(None);
833 track.finish().unwrap();
834
835 let mut consumer = Consumer::<Value>::new(consumer_track, ConsumerConfig::default());
836 let waiter = kio::Waiter::noop();
837
838 assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
840
841 group
842 .write_frame(
843 moq_net::Timestamp::ZERO,
844 Bytes::from(serde_json::to_vec(&json!({ "a": 1 })).unwrap()),
845 )
846 .unwrap();
847 group.finish().unwrap();
848
849 match consumer.poll_next(&waiter) {
850 Poll::Ready(Ok(Some(value))) => assert_eq!(value, json!({ "a": 1 })),
851 other => panic!("expected the catalog value, got {other:?}"),
852 }
853 }
854
855 #[test]
856 fn late_joiner_collapses_backlog_to_latest() {
857 let (mut producer, track) = producer(cfg(100));
860 for n in 0..=20 {
861 producer.update(&json!({ "n": n })).unwrap();
862 }
863 producer.finish().unwrap();
864
865 assert_eq!(track.latest(), Some(0));
867 let values = drain(track);
868 assert_eq!(
869 values,
870 vec![json!({ "n": 20 })],
871 "backlog should collapse to the latest value"
872 );
873 }
874
875 #[test]
876 fn compressed_late_joiner_collapses_backlog_to_latest() {
877 let (mut producer, track) = producer(cfg_deflate(100));
879 for n in 0..=20 {
880 producer.update(&json!({ "n": n })).unwrap();
881 }
882 producer.finish().unwrap();
883
884 assert_eq!(track.latest(), Some(0));
885 let values = drain_with(deflate_consumer(track));
886 assert_eq!(
887 values,
888 vec![json!({ "n": 20 })],
889 "compressed backlog should collapse to the latest"
890 );
891 }
892
893 #[test]
894 fn compressed_snapshot_per_group_roundtrips() {
895 let (mut producer, track) = producer(cfg_deflate(0));
896 producer.update(&json!({ "a": 1 })).unwrap();
897 producer.update(&json!({ "a": 2 })).unwrap();
898 producer.finish().unwrap();
899
900 assert_eq!(track.latest(), Some(1));
902 let values = drain_with(deflate_consumer(track));
903 assert_eq!(values, vec![json!({ "a": 2 })]);
904 }
905
906 #[test]
907 fn compressed_deltas_share_one_group() {
908 let (mut producer, track) = producer(cfg_deflate(100));
909 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
910 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
911 producer.update(&json!({ "a": 1, "b": 3 })).unwrap();
912 producer.finish().unwrap();
913
914 assert_eq!(track.latest(), Some(0));
916 let values = drain_with(deflate_consumer(track));
917 assert_eq!(values.last().unwrap(), &json!({ "a": 1, "b": 3 }));
918 }
919
920 #[test]
921 fn compressed_late_joiner_reconstructs_from_deltas() {
922 let (mut producer, track) = producer(cfg_deflate(100));
923 producer.update(&json!({ "a": 1, "b": 1 })).unwrap();
924 producer.update(&json!({ "a": 1, "b": 2 })).unwrap();
925 producer.update(&json!({ "a": 5, "b": 2 })).unwrap();
926 producer.finish().unwrap();
927
928 let values = drain_with(deflate_consumer(track));
930 assert_eq!(values.last().unwrap(), &json!({ "a": 5, "b": 2 }));
931 }
932
933 #[test]
934 fn compressed_deltas_roll_on_compressed_budget() {
935 let (mut producer, track) = producer(cfg_deflate(2));
941 for n in 0..=40 {
942 producer.update(&json!({ "n": n })).unwrap();
943 }
944 producer.finish().unwrap();
945
946 assert!(
947 track.latest().unwrap() > 0,
948 "a tight ratio should roll at least one compressed group"
949 );
950 assert_eq!(drain_with(deflate_consumer(track)).last().unwrap(), &json!({ "n": 40 }));
951 }
952
953 #[test]
954 fn compression_shrinks_wire_frames() {
955 let value = json!({ "renditions": ["video".repeat(50), "video".repeat(50), "video".repeat(50)] });
957
958 let plaintext_bytes = wire_frame_len(cfg(0), &value);
959 let compressed_bytes = wire_frame_len(cfg_deflate(0), &value);
960 assert!(
961 compressed_bytes < plaintext_bytes,
962 "compressed frame {compressed_bytes} should be smaller than plaintext {plaintext_bytes}"
963 );
964 }
965
966 #[test]
967 fn compressed_deltas_reuse_window() {
968 let (mut producer, mut track) = producer(cfg_deflate(100));
971 let phrase = "Media over QUIC delivers real-time latency at massive scale";
972 producer.update(&json!({ "note": phrase })).unwrap();
973 producer.update(&json!({ "note": phrase, "echo": phrase })).unwrap();
974 producer.finish().unwrap();
975
976 let waiter = kio::Waiter::noop();
978 let Poll::Ready(Ok(Some(mut group))) = track.poll_next_group(&waiter) else {
979 panic!("expected a group");
980 };
981 let mut frames = Vec::new();
982 while let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) {
983 frames.push(frame.payload);
984 }
985 assert_eq!(frames.len(), 2, "snapshot + one delta in a single group");
986
987 let raw_delta = serde_json::to_vec(&json!({ "echo": phrase })).unwrap();
989 assert!(
990 frames[1].len() < raw_delta.len() / 2,
991 "windowed delta {} should be far below the raw patch {}",
992 frames[1].len(),
993 raw_delta.len()
994 );
995 }
996
997 fn wire_frame_len(config: ProducerConfig, value: &Value) -> usize {
999 let (mut producer, mut track) = producer(config);
1000 producer.update(value).unwrap();
1001 producer.finish().unwrap();
1002
1003 let waiter = kio::Waiter::noop();
1004 let Poll::Ready(Ok(Some(mut group))) = track.poll_next_group(&waiter) else {
1005 panic!("expected a group");
1006 };
1007 let Poll::Ready(Ok(Some(frame))) = group.poll_read_frame(&waiter) else {
1009 panic!("expected a frame");
1010 };
1011 frame.payload.len()
1012 }
1013}