Skip to main content

deser_core/de/
stream.rs

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