Skip to main content

deser_core/de/
stream.rs

1use crate::de::DeserializeDriver;
2use crate::error::{Error, ErrorKind};
3
4/// The result of [`StreamDeserializer::frame`].
5#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
6pub enum Frame {
7    /// The next value is complete.
8    ///
9    /// Its bytes are `input[start..end]`.  Afterwards the first `consumed`
10    /// bytes of the input (at least up to `end`) are discarded, the bytes
11    /// before `start` are skipped.
12    Value {
13        start: usize,
14        end: usize,
15        consumed: usize,
16    },
17    /// The input does not contain a complete value.
18    ///
19    /// The first `consumed` bytes of the input are discarded, for instance
20    /// whitespace before the next value.  If bytes were consumed, the
21    /// deserializer is invoked again right away (as a value might follow
22    /// them), otherwise once more input was read.  At the end of the input
23    /// the deserializer must not return this without consuming bytes.
24    Incomplete { consumed: usize },
25    /// There are no more values.
26    ///
27    /// This must only be returned at the end of the input.
28    End,
29}
30
31/// The result of [`StreamDeserializer::feed`].
32#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
33pub enum Progress {
34    /// The value is complete, it used the first `consumed` bytes of the
35    /// input.
36    Done { consumed: usize },
37    /// The value needs more input.  The first `consumed` bytes of the input
38    /// were used and are discarded, the next call continues with the input
39    /// after them (followed by the new data).
40    NeedMore { consumed: usize },
41    /// There are no more values.
42    ///
43    /// This must only be returned at the end of the input.
44    End,
45}
46
47/// Deserializes a stream of values from input that arrives in chunks.
48///
49/// This is implemented by the stream deserializers of the data formats
50/// (for instance `deser_json::StreamDeserializer`).  Unlike a
51/// [`Deserializer`](crate::de::Deserializer), which pulls values from an
52/// input it holds (and can lend data from it), a stream deserializer is
53/// given the input as it arrives.  It holds everything a stream needs to
54/// remember: the progress of the scan and what earlier parts of the stream
55/// established for the values that follow, for instance the names of the
56/// columns of a CSV file.
57///
58/// Stream deserializers do not do IO: a
59/// [`stream::InputBuffer`](crate::stream::InputBuffer) holds the input and
60/// invokes them, the readers of `deser::io` and of other IO adapters
61/// (like `deser-tokio`) fill the buffer.
62///
63/// # Frames
64///
65/// A stream deserializer splits the input into frames: it finds the bytes
66/// of the next value in the input that was read so far (see
67/// [`frame`](Self::frame)), for instance a line with JSON Lines.  Once a
68/// value is complete it's deserialized from its frame (see
69/// [`drive_frame`](Self::drive_frame)), typically with the format's
70/// regular parser.  Values read from their frames can borrow from the
71/// buffer.
72///
73/// # Feeding
74///
75/// Formats which can be parsed while the input arrives (like JSON and
76/// CBOR) can also deserialize values while their input is fed to them
77/// (see [`feed`](Self::feed)).  Only incomplete tokens are buffered, so
78/// the memory used does not depend on the size of the values.  Values
79/// read this way cannot borrow from the input.
80///
81/// ```
82/// use deser::de::{DeserializeDriver, Frame, StreamDeserializer};
83/// use deser::stream::{InputBuffer, Status};
84/// use deser::Error;
85///
86/// /// A format with a number per line.
87/// struct Lines;
88///
89/// impl StreamDeserializer for Lines {
90///     fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
91///         Ok(match input.iter().position(|&b| b == b'\n') {
92///             Some(end) => Frame::Value { start: 0, end, consumed: end + 1 },
93///             None if eof && input.is_empty() => Frame::End,
94///             None if eof => Frame::Value { start: 0, end: input.len(), consumed: input.len() },
95///             None => Frame::Incomplete { consumed: 0 },
96///         })
97///     }
98///
99///     fn drive_frame<'de>(
100///         &mut self,
101///         frame: &'de [u8],
102///         driver: &mut DeserializeDriver<'_, 'de>,
103///     ) -> Result<(), Error> {
104///         let value: u64 = std::str::from_utf8(frame).unwrap().parse().unwrap();
105///         driver.emit(value)
106///     }
107/// }
108///
109/// let mut buffer = InputBuffer::new(Lines);
110/// buffer.extend_from_slice(b"1\n2");
111/// assert_eq!(buffer.poll().unwrap(), Status::Ready);
112/// assert_eq!(buffer.deserialize::<u32>().unwrap(), 1);
113/// assert_eq!(buffer.poll().unwrap(), Status::NeedInput);
114/// buffer.set_eof();
115/// assert_eq!(buffer.poll().unwrap(), Status::Ready);
116/// assert_eq!(buffer.deserialize::<u32>().unwrap(), 2);
117/// assert_eq!(buffer.poll().unwrap(), Status::End);
118/// ```
119pub trait StreamDeserializer {
120    /// Finds the next value in the input.
121    ///
122    /// The input holds the data that was read so far (minus the data that
123    /// was discarded).  If it does not contain a complete value yet,
124    /// [`Frame::Incomplete`] is returned and the method is invoked again
125    /// once more data was read: the input then starts after the bytes that
126    /// were consumed and continues with the new data.  This allows
127    /// deserializers to keep the progress of their scan so they do not have
128    /// to scan the input again.  `eof` is `true` if no more data follows
129    /// the input.
130    ///
131    /// Once a value is complete, [`Frame::Value`] is returned and the value
132    /// is deserialized with [`drive_frame`](Self::drive_frame).  The next
133    /// call starts a new value, again after the consumed bytes.  Offsets
134    /// of errors refer to the input.
135    fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error>;
136
137    /// Deserializes a value from its frame.
138    ///
139    /// The frame holds the bytes of a value found by
140    /// [`frame`](Self::frame), it's deserialized right after it was found.
141    /// This allows deserializers to keep what they learned while scanning
142    /// the frame (like the positions of fields) so they do not have to scan
143    /// it again.  Events can borrow from the frame.  Offsets of errors
144    /// refer to the frame.
145    ///
146    /// Data that is only valid for the call (for instance names kept by
147    /// the deserializer) can be emitted without copying it with
148    /// [`DeserializeDriver::emit`].
149    fn drive_frame<'de>(
150        &mut self,
151        frame: &'de [u8],
152        driver: &mut DeserializeDriver<'_, 'de>,
153    ) -> Result<(), Error>;
154
155    /// Returns `true` if the format is text.
156    ///
157    /// For text formats the positions of errors are resolved into lines and
158    /// columns.  This is `false` by default.
159    fn is_text(&self) -> bool {
160        false
161    }
162
163    /// Returns `true` if the deserializer implements [`feed`](Self::feed).
164    ///
165    /// This can depend on the configuration, for instance JSON Lines are
166    /// read line by line.
167    fn supports_feed(&self) -> bool {
168        false
169    }
170
171    /// Deserializes a value while its input arrives.
172    ///
173    /// This is only invoked if [`supports_feed`](Self::supports_feed)
174    /// returns `true`.  It's used instead of [`frame`](Self::frame) and
175    /// [`drive_frame`](Self::drive_frame) for values which do not borrow
176    /// from the input.  The deserializer emits the events of the parts of
177    /// the value in the input into the driver and returns how much of the
178    /// input it used ([`Progress::NeedMore`]) until the value is complete
179    /// ([`Progress::Done`]).  Only incomplete tokens need to be kept.  The
180    /// driver is the same for all calls for a value, the first call for a
181    /// value starts where the previous value ended.  At the end of the
182    /// input (`eof`) the value has to be completed (or fail).
183    ///
184    /// As the input does not live beyond the call, the events cannot
185    /// borrow from it.  `offset` is the offset of the input in the stream:
186    /// the input ranges of the events and the offsets of errors refer to
187    /// positions in the stream.  After an error the value is abandoned, the
188    /// deserializer decides if the stream can continue with the next value
189    /// (for instance by skipping the rest of the value if a sink failed) or
190    /// if further calls fail.
191    fn feed(
192        &mut self,
193        input: &[u8],
194        offset: usize,
195        eof: bool,
196        driver: &mut DeserializeDriver<'_, '_>,
197    ) -> Result<Progress, Error> {
198        let _ = (input, offset, eof, driver);
199        Err(Error::new(
200            ErrorKind::Unexpected,
201            "the deserializer cannot deserialize while the input arrives",
202        ))
203    }
204
205    /// Finds the start of the next value without deserializing it.
206    ///
207    /// This is used to check if another value follows (see
208    /// [`InputBuffer::peek`](crate::stream::InputBuffer::peek)) without
209    /// reading the value.  It skips what precedes the next value (for
210    /// instance whitespace) and returns:
211    ///
212    /// * `Some(Progress::Done { consumed })` if a value starts after the
213    ///   first `consumed` bytes of the input (which are discarded).
214    /// * `Some(Progress::NeedMore { consumed })` if more input is needed
215    ///   to know, the first `consumed` bytes are discarded.
216    /// * `Some(Progress::End)` if there are no more values.  This must only
217    ///   be returned at the end of the input.
218    ///
219    /// Offsets of errors refer to the input.  The provided implementation
220    /// returns `None`, then the next value is found by framing it (which
221    /// buffers it completely).  Deserializers that support
222    /// [`feed`](Self::feed) should implement this so values are not
223    /// buffered to find out if they exist.
224    fn peek(&mut self, input: &[u8], eof: bool) -> Result<Option<Progress>, Error> {
225        let _ = (input, eof);
226        Ok(None)
227    }
228}
229
230impl<D: StreamDeserializer + ?Sized> StreamDeserializer for &mut D {
231    fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
232        (**self).frame(input, eof)
233    }
234
235    fn drive_frame<'de>(
236        &mut self,
237        frame: &'de [u8],
238        driver: &mut DeserializeDriver<'_, 'de>,
239    ) -> Result<(), Error> {
240        (**self).drive_frame(frame, driver)
241    }
242
243    fn is_text(&self) -> bool {
244        (**self).is_text()
245    }
246
247    fn supports_feed(&self) -> bool {
248        (**self).supports_feed()
249    }
250
251    fn feed(
252        &mut self,
253        input: &[u8],
254        offset: usize,
255        eof: bool,
256        driver: &mut DeserializeDriver<'_, '_>,
257    ) -> Result<Progress, Error> {
258        (**self).feed(input, offset, eof, driver)
259    }
260
261    fn peek(&mut self, input: &[u8], eof: bool) -> Result<Option<Progress>, Error> {
262        (**self).peek(input, eof)
263    }
264}