moq_json/snapshot/consumer.rs
1//! Consuming a JSON value from a track: a [`Decoder`] plus the track it reads from.
2
3use std::task::Poll;
4
5use serde::de::DeserializeOwned;
6
7use super::Decoder;
8use crate::{Compression, Result};
9
10/// Track-owning options for a [`Consumer`].
11///
12/// Build from [`Default`] and override fields (the struct is `#[non_exhaustive]`, so new options
13/// stay additive).
14#[derive(Debug, Clone, Default)]
15#[non_exhaustive]
16pub struct Config {
17 /// How the frames are compressed. Must match the encoder's
18 /// [`Config::compression`](super::Config::compression). Defaults to [`Compression::None`].
19 pub compression: Compression,
20
21 /// Maximum bytes in an inflated frame or the compact JSON of the reconstructed value.
22 /// Unset preserves the DEFLATE decoder's default frame cap and leaves state size unlimited.
23 pub max_size: Option<usize>,
24}
25
26/// Consumes a JSON value from a track, reconstructing it from snapshots and deltas.
27///
28/// A [`Decoder`] that owns its track: it reads groups, routes each frame by its position, and
29/// yields the reconstructed value. When something else already owns the track, use the [`Decoder`]
30/// directly.
31pub struct Consumer<T> {
32 track: moq_net::track::Ordered,
33 group: Option<moq_net::group::Consumer>,
34 decoder: Decoder<T>,
35 frames_read: usize,
36}
37
38impl<T: DeserializeOwned> Consumer<T> {
39 /// Create a consumer reading from the given track subscriber.
40 ///
41 /// Set [`Config::compression`] to read a track written by a producer with the same
42 /// [`compression`](super::Config::compression).
43 pub fn new(track: moq_net::track::Subscriber, config: Config) -> Self {
44 Self {
45 track: track.ordered(),
46 group: None,
47 decoder: Decoder::new(config),
48 frames_read: 0,
49 }
50 }
51
52 /// Get the next reconstructed value, or `None` once the track ends.
53 pub async fn next(&mut self) -> Result<Option<T>>
54 where
55 T: Unpin,
56 {
57 kio::wait(|waiter| self.poll_next(waiter)).await
58 }
59
60 /// Poll for the next reconstructed value, without blocking.
61 ///
62 /// Jumps to the newest group, reads its snapshot, and applies deltas in order. All frames already
63 /// buffered in the group are applied in one poll but only the resulting *latest* value is yielded:
64 /// the intermediate reconstructions are stale, so a late joiner (or any consumer that has fallen
65 /// behind) catches up to the head in a single step instead of replaying every superseded state.
66 /// Frames must still be decoded in order (the DEFLATE window and merge patches are sequential);
67 /// only the per-frame deserialize and yield are skipped. Switching to a newer group discards the
68 /// older one.
69 ///
70 /// A group the transport can no longer serve is discarded the same way, not reported: on a
71 /// snapshot track its content is superseded by definition, so the reader waits for the
72 /// replacement. Only a failure of the track itself ends the stream.
73 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<T>>> {
74 // Drain to the newest group, resetting reconstruction state whenever we switch.
75 let track_finished = loop {
76 match self.track.poll_next_group(waiter)? {
77 Poll::Ready(Some(group)) => {
78 self.group = Some(group);
79 // The next frame is the new group's snapshot, which also restarts the decoder's window.
80 self.frames_read = 0;
81 }
82 Poll::Ready(None) => break true,
83 Poll::Pending => break false,
84 }
85 };
86
87 // Apply every frame currently buffered in the group, tracking whether any moved us forward and
88 // whether the group is still open with nothing buffered yet (vs. exhausted).
89 // `poll_read_frame` returns an owned `Poll`, so the borrow of `self.group` ends before the
90 // match arms, leaving `apply` (and clearing the group) free to take `&mut self`.
91 let mut advanced = false;
92 let mut group_pending = false;
93 while let Some(group) = &mut self.group {
94 match group.poll_read_frame(waiter) {
95 Poll::Ready(Ok(Some(frame))) => {
96 self.apply(&frame.payload)?;
97 advanced = true;
98 }
99 // The current group is exhausted; wait for a newer one.
100 Poll::Ready(Ok(None)) => {
101 self.group = None;
102 break;
103 }
104 // The transport can no longer serve the rest of this group: it was superseded and
105 // reclaimed (`Old`), dropped under memory pressure (`Evicted`), or read past the
106 // drift budget (`Lagged`). A snapshot reader only ever wants the newest value, so a
107 // group whose content is gone is never fatal: drop it and wait for its replacement.
108 // A track- or session-level failure still arrives through `poll_next_group` above.
109 Poll::Ready(Err(err)) => {
110 let sequence = group.sequence;
111 self.group = None;
112 tracing::warn!(
113 track = self.track.name(),
114 group = sequence,
115 error = ?err,
116 "snapshot group lost; waiting for a newer one"
117 );
118 break;
119 }
120 // The group is still open but has nothing buffered yet.
121 Poll::Pending => {
122 group_pending = true;
123 break;
124 }
125 }
126 }
127
128 if advanced {
129 // Deserialize once, from the head of the backlog we just drained.
130 return Poll::Ready(Ok(self.decoder.decode()?));
131 }
132
133 // An open group may still deliver frames even after the track finishes (it was appended before
134 // the finish), so wait on it rather than ending the stream.
135 if group_pending {
136 return Poll::Pending;
137 }
138
139 if track_finished {
140 Poll::Ready(Ok(None))
141 } else {
142 Poll::Pending
143 }
144 }
145
146 /// Apply one frame: frame 0 of a group is a snapshot, the rest are merge patches.
147 fn apply(&mut self, payload: &[u8]) -> Result<()> {
148 match self.frames_read {
149 0 => self.decoder.snapshot(payload)?,
150 _ => self.decoder.delta(payload)?,
151 }
152 self.frames_read += 1;
153 Ok(())
154 }
155}