Skip to main content

deser_jsonc/
stream.rs

1// @generated from deser-template-json/src/stream.rs by
2// deser-template-json/generate.py.  Do not edit.
3//! Reading JSON streams.
4#[cfg(feature = "io")]
5use std::io::Read;
6
7#[cfg(feature = "io")]
8use deser_core::de::DeserializeOwned;
9use deser_core::de::{self, DeserializeDriver, Frame, Progress};
10use deser_core::{Error, ErrorKind, State};
11
12use crate::Trailing;
13use crate::de::{Deserializer, DeserializerConfig};
14use crate::parser::Cursor;
15use crate::parser::{Copying, Discard, Options, Parser, Progress as ParseProgress};
16use crate::scan::LineScan;
17use crate::scan::skip_to_escape;
18
19/// Returns where the next token starts and if it's there.
20///
21/// If the input ends within a comment (and more input follows), this is
22/// where the comment starts so that it's scanned again with more input,
23/// the token is not there.
24fn skip_whitespace(input: &[u8], pos: usize, eof: bool) -> (usize, bool) {
25    let mut cursor = Cursor::new_partial(input, pos, eof);
26    let token = cursor.parse_whitespace().is_some();
27    (cursor.pos, token)
28}
29
30/// The state of a JSON stream that is read.
31#[derive(Debug, Default)]
32struct StreamState {
33    // `Trailing::Strict`: the value was read
34    done: bool,
35    // parses values while their input arrives (see `drive_partial`)
36    parser: Parser,
37    // the rest of a value that failed in a sink is skipped from the
38    // position in the input
39    skipping: Option<usize>,
40    // parsing failed, the stream cannot be continued
41    failed: bool,
42    // the stream ended within a value
43    ended: bool,
44    // the position up to which the input was scanned
45    pos: usize,
46    // `Trailing::Newline`: the scan of the current line
47    line: LineScan,
48    // `Trailing::Stop`: the value being scanned
49    value: Option<Value>,
50}
51
52/// The state of the scan of a value.
53#[derive(Debug)]
54struct Value {
55    start: usize,
56    kind: ValueKind,
57    depth: usize,
58    in_string: bool,
59}
60
61#[derive(Debug)]
62enum ValueKind {
63    /// A number or literal, which is scanned by parsing it (`1-2` is two
64    /// numbers).  The parser consumed the input up to `start + parsed`.
65    Scalar { parser: Parser, parsed: usize },
66    /// A string, map or sequence.
67    Structure,
68}
69
70fn frame_all(state: &mut StreamState, input: &[u8], eof: bool) -> Result<Frame, Error> {
71    if state.done {
72        // only whitespace may follow the value
73        return trailing_whitespace(input, 0, eof).map(|progress| match progress {
74            Progress::End => Frame::End,
75            // an incomplete comment is scanned again with more input
76            Progress::NeedMore { consumed } => Frame::Incomplete { consumed },
77            Progress::Done { .. } => unreachable!(),
78        });
79    }
80    if !eof {
81        return Ok(Frame::Incomplete { consumed: 0 });
82    }
83    // an unterminated comment is reported by the parser
84    let (start, _) = skip_whitespace(input, 0, eof);
85    Ok(if start < input.len() {
86        state.done = true;
87        Frame::Value {
88            start,
89            end: input.len(),
90            consumed: input.len(),
91        }
92    } else {
93        Frame::End
94    })
95}
96
97/// Checks that only whitespace follows a value.
98///
99/// `offset` is the offset of the input in the stream for errors.
100fn trailing_whitespace(input: &[u8], offset: usize, eof: bool) -> Result<Progress, Error> {
101    match skip_whitespace(input, 0, eof) {
102        (pos, _) if pos == input.len() && eof => Ok(Progress::End),
103        // an incomplete comment is scanned again with more input
104        (pos, false) if !eof => Ok(Progress::NeedMore { consumed: pos }),
105        (pos, _) => Err(Error::with_offset(
106            ErrorKind::Syntax,
107            "garbage after input",
108            offset + pos,
109        )),
110    }
111}
112
113fn frame_line(state: &mut StreamState, input: &[u8], eof: bool) -> Frame {
114    // whitespace before the value is skipped first like the deserializer
115    // does, comments in it can span lines
116    if state.pos == 0 {
117        match skip_whitespace(input, 0, eof) {
118            (0, true) => {}
119            (start, true) => return Frame::Incomplete { consumed: start },
120            (_, false) if input.is_empty() && eof => return Frame::End,
121            // an unterminated comment is reported by the parser
122            (start, false) if eof && start < input.len() => {
123                return Frame::Value {
124                    start,
125                    end: input.len(),
126                    consumed: input.len(),
127                };
128            }
129            // an incomplete comment is scanned again with more input
130            (start, false) => return Frame::Incomplete { consumed: start },
131        }
132    }
133    // line breaks in comments and strings do not end the line
134    let end = state.line.find_end(input, state.pos);
135    let end = match end {
136        Some(end) => end,
137        None if eof => input.len(),
138        None => {
139            state.pos = input.len();
140            return Frame::Incomplete { consumed: 0 };
141        }
142    };
143    state.pos = 0;
144    state.line = LineScan::default();
145    let consumed = (end + 1).min(input.len());
146    match skip_whitespace(&input[..end], 0, true) {
147        (start, _) if start < end => Frame::Value {
148            start,
149            end,
150            consumed,
151        },
152        _ if consumed == 0 => Frame::End,
153        // blank lines are skipped
154        _ => Frame::Incomplete { consumed },
155    }
156}
157
158fn frame_value(state: &mut StreamState, input: &[u8], eof: bool) -> Frame {
159    let value = match state.value {
160        Some(ref mut value) => value,
161        None => {
162            let start = match skip_whitespace(input, 0, eof) {
163                (start, true) => start,
164                (_, false) if input.is_empty() && eof => return Frame::End,
165                (consumed, false) => return Frame::Incomplete { consumed },
166            };
167            let (kind, depth, in_string) = match input[start] {
168                b'"' => (ValueKind::Structure, 0, true),
169                b'{' | b'[' => (ValueKind::Structure, 1, false),
170                // a value cannot start with these, the parser reports the
171                // error
172                b'}' | b']' | b',' | b':' => {
173                    return Frame::Value {
174                        start,
175                        end: start + 1,
176                        consumed: start + 1,
177                    };
178                }
179                _ => (
180                    ValueKind::Scalar {
181                        parser: Parser::default(),
182                        parsed: 0,
183                    },
184                    0,
185                    false,
186                ),
187            };
188            state.pos = start + 1;
189            state.value.insert(Value {
190                start,
191                kind,
192                depth,
193                in_string,
194            })
195        }
196    };
197
198    let end = match value.kind {
199        ValueKind::Scalar {
200            ref mut parser,
201            ref mut parsed,
202        } => {
203            let pos = value.start + *parsed;
204            let options = Options {
205                validate_utf8: true,
206                exact_numbers: false,
207            };
208            let mut discard = Discard(State::new());
209            match parser.parse(&input[pos..], 0, eof, 0, options, &mut discard) {
210                Ok(ParseProgress::Done(end)) => Some(pos + end),
211                Ok(ParseProgress::NeedMore(consumed)) => {
212                    *parsed += consumed;
213                    state.pos = input.len();
214                    None
215                }
216                // the value ends at the error, the parser reports it when
217                // the value is deserialized
218                Err(err) => Some(error_end(&err, input, pos)),
219            }
220        }
221        ValueKind::Structure => scan_structure(input, &mut state.pos, value),
222    };
223
224    let start = value.start;
225    match end {
226        Some(end) => {
227            state.value = None;
228            state.pos = 0;
229            Frame::Value {
230                start,
231                end,
232                consumed: end,
233            }
234        }
235        // the value is incomplete, the parser reports the error
236        None if eof => {
237            state.value = None;
238            state.pos = 0;
239            Frame::Value {
240                start,
241                end: input.len(),
242                consumed: input.len(),
243            }
244        }
245        None => {
246            // discard the whitespace before the value
247            value.start = 0;
248            state.pos -= start;
249            Frame::Incomplete { consumed: start }
250        }
251    }
252}
253
254/// Returns where the frame of a value which failed to parse ends.
255///
256/// The frame ends after the byte that failed, the parser fails there
257/// again when the value is deserialized.  If the input ended within the
258/// value (where the error is located at the start of the token, like a
259/// comment, which is valid if it's cut off) the frame is the rest of the
260/// input.
261fn error_end(err: &Error, input: &[u8], pos: usize) -> usize {
262    match err.offset() {
263        Some(offset) if err.kind() != ErrorKind::EndOfFile => (pos + offset + 1).min(input.len()),
264        _ => input.len(),
265    }
266}
267
268/// Scans a string, map or sequence from `pos`.
269///
270/// Returns the end of the value if it's complete.
271fn scan_structure(input: &[u8], pos: &mut usize, value: &mut Value) -> Option<usize> {
272    let mut index = *pos;
273    while index < input.len() {
274        if value.in_string {
275            index = skip_to_escape(input, index);
276            let byte = input.get(index).copied();
277            match byte {
278                Some(b'"') => {
279                    value.in_string = false;
280                    index += 1;
281                    if value.depth == 0 {
282                        return Some(index);
283                    }
284                }
285                Some(b'\\') => {
286                    // the escaped byte is scanned again with more input
287                    if index + 1 >= input.len() {
288                        break;
289                    }
290                    index += 2;
291                }
292                // control characters are reported by the parser
293                Some(_) => index += 1,
294                None => break,
295            }
296        } else {
297            match input[index] {
298                b'"' => {
299                    value.in_string = true;
300                }
301                // comments are skipped, incomplete ones are scanned again
302                // with more input
303                b'/' => {
304                    match skip_whitespace(input, index, false) {
305                        // not a comment, the parser reports the error
306                        (next, true) if next == index => index += 1,
307                        (next, token) => {
308                            index = next;
309                            if !token && next < input.len() {
310                                break;
311                            }
312                        }
313                    }
314                    continue;
315                }
316                b'{' | b'[' => value.depth += 1,
317                b'}' | b']' => {
318                    value.depth -= 1;
319                    if value.depth == 0 {
320                        return Some(index + 1);
321                    }
322                }
323                _ => {}
324            }
325            index += 1;
326        }
327    }
328    *pos = index;
329    None
330}
331
332/// Reads a stream of JSON values (see [`deser::stream`](deser_core::stream)).
333///
334/// How the stream is split depends on [`DeserializerConfig::set_trailing`]:
335///
336/// * [`Trailing::Strict`]: the stream holds a single value which is parsed
337///   once the whole stream was read.
338/// * [`Trailing::Newline`]: every line holds a value ([JSON
339///   Lines](https://jsonlines.org/)).  Blank lines are skipped and errors
340///   only discard their line, reading continues with the next one.
341/// * [`Trailing::Stop`]: values follow each other (optionally separated by
342///   whitespace) and are split where they end.  Numbers at the end of the
343///   stream are only complete at the end of the stream, other values are
344///   complete once their last byte was read.  Reading continues after
345///   values that fail to deserialize, values that are not valid JSON end
346///   the stream (unless they are read from their frames).
347///
348/// Except for JSON Lines, values which do not borrow are deserialized while
349/// their input arrives (see
350/// [`StreamDeserializer::drive_partial`](de::StreamDeserializer::drive_partial))
351/// so only incomplete tokens are buffered.
352///
353/// ```
354/// # #[cfg(feature = "io")] {
355/// use deser_jsonc::{DeserializerConfig, Trailing};
356///
357/// const LINES: DeserializerConfig =
358///     DeserializerConfig::builder().trailing(Trailing::Newline).build();
359/// let mut reader = LINES.reader(&b"[1, 2]\n[3]\n"[..]);
360/// assert_eq!(reader.read::<Vec<u32>>().unwrap(), Some(vec![1, 2]));
361/// assert_eq!(reader.read::<Vec<u32>>().unwrap(), Some(vec![3]));
362/// assert_eq!(reader.read::<Vec<u32>>().unwrap(), None);
363/// # }
364/// ```
365///
366/// Values which are read from their frames are parsed like with a
367/// [`Deserializer`], so they can borrow from the stream's buffer (see
368/// [`InputBuffer::deserialize`](deser_core::stream::InputBuffer::deserialize)).  The input ranges (and thus
369/// locations) of these values refer to the start of their line (or value),
370/// those of values that are deserialized while their input arrives to the
371/// stream.
372#[derive(Debug)]
373pub struct StreamDeserializer {
374    config: DeserializerConfig,
375    state: StreamState,
376}
377
378impl Default for StreamDeserializer {
379    fn default() -> StreamDeserializer {
380        StreamDeserializer::new()
381    }
382}
383
384impl StreamDeserializer {
385    /// Creates a stream deserializer.
386    pub fn new() -> StreamDeserializer {
387        StreamDeserializer::with_config(DeserializerConfig::new())
388    }
389
390    /// Creates a stream deserializer with the given configuration.
391    pub fn with_config(config: DeserializerConfig) -> StreamDeserializer {
392        StreamDeserializer {
393            config,
394            state: StreamState::default(),
395        }
396    }
397
398    /// Returns the configuration.
399    pub fn config(&self) -> &DeserializerConfig {
400        &self.config
401    }
402
403    /// Skips what precedes the next value.
404    ///
405    /// This skips the rest of a value that failed in a sink and the
406    /// whitespace before the next value.  Returns where the value starts,
407    /// or the progress if there is no value (yet).
408    fn skip_to_value(
409        &mut self,
410        input: &[u8],
411        offset: usize,
412        eof: bool,
413    ) -> Result<Result<usize, Progress>, Error> {
414        let options = self.options();
415        let state = &mut self.state;
416        if state.ended {
417            return Ok(Err(Progress::End));
418        }
419        if state.failed {
420            return Err(Error::new(
421                ErrorKind::InvalidState,
422                "cannot continue after an error",
423            ));
424        }
425
426        // skip the rest of a value that failed in a sink
427        let mut pos = 0;
428        if let Some(skip) = state.skipping {
429            let mut discard = Discard(State::new());
430            match state
431                .parser
432                .parse(input, skip, eof, offset, options, &mut discard)
433            {
434                Ok(ParseProgress::Done(end)) => {
435                    state.skipping = None;
436                    state.done = self.config.trailing_mode() == Trailing::Strict;
437                    pos = end;
438                }
439                Ok(ParseProgress::NeedMore(consumed)) => {
440                    state.skipping = Some(0);
441                    return Ok(Err(Progress::NeedMore { consumed }));
442                }
443                Err(err) => {
444                    state.parser.reset();
445                    state.skipping = None;
446                    state.failed = true;
447                    return Err(err);
448                }
449            }
450        }
451
452        if !state.parser.is_idle() {
453            return Ok(Ok(pos));
454        }
455        if state.done {
456            return match trailing_whitespace(&input[pos..], offset + pos, eof)? {
457                Progress::NeedMore { consumed } => Ok(Err(Progress::NeedMore {
458                    consumed: pos + consumed,
459                })),
460                progress => Ok(Err(progress)),
461            };
462        }
463        // a new value, skip the whitespace before it
464        // (an incomplete comment is scanned again with more input)
465        let token;
466        (pos, token) = skip_whitespace(input, pos, eof);
467        if !token && (pos == input.len() || !eof) {
468            return Ok(Err(if eof {
469                Progress::End
470            } else {
471                Progress::NeedMore { consumed: pos }
472            }));
473        }
474        Ok(Ok(pos))
475    }
476
477    /// Returns the options of the parser.
478    fn options(&self) -> Options {
479        Options {
480            validate_utf8: true,
481            exact_numbers: self.config.exact_numbers_enabled(),
482        }
483    }
484}
485
486impl de::StreamDeserializer for StreamDeserializer {
487    fn context(&self) -> deser_core::Context {
488        self.config.context().clone()
489    }
490
491    fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
492        let state = &mut self.state;
493        match self.config.trailing_mode() {
494            Trailing::Strict => frame_all(state, input, eof),
495            Trailing::Newline => Ok(frame_line(state, input, eof)),
496            Trailing::Stop => Ok(frame_value(state, input, eof)),
497        }
498    }
499
500    fn drive_frame<'de>(
501        &mut self,
502        frame: &'de [u8],
503        driver: &mut DeserializeDriver<'_, 'de>,
504    ) -> Result<(), Error> {
505        let mut de = Deserializer::from_frame(frame, &self.config);
506        de.drive(driver)
507    }
508
509    fn is_text(&self) -> bool {
510        true
511    }
512
513    /// JSON Lines are read line by line, the other values while their
514    /// input arrives.
515    fn supports_partial(&self) -> bool {
516        self.config.trailing_mode() != Trailing::Newline
517    }
518
519    fn drive_partial(
520        &mut self,
521        input: &[u8],
522        offset: usize,
523        eof: bool,
524        driver: &mut DeserializeDriver<'_, '_>,
525    ) -> Result<Progress, Error> {
526        let pos = match self.skip_to_value(input, offset, eof)? {
527            Ok(pos) => pos,
528            Err(progress) => return Ok(progress),
529        };
530        let options = self.options();
531        let state = &mut self.state;
532        match state
533            .parser
534            .parse(input, pos, eof, offset, options, &mut Copying(driver))
535        {
536            Ok(ParseProgress::Done(end)) => {
537                state.done = self.config.trailing_mode() == Trailing::Strict;
538                Ok(Progress::Done { consumed: end })
539            }
540            Ok(ParseProgress::NeedMore(consumed)) => Ok(Progress::NeedMore { consumed }),
541            Err(err) => {
542                if let Some(resume) = state.parser.recoverable() {
543                    // a sink failed, the next call continues after the
544                    // value.  The input is not consumed on errors.
545                    state.skipping = Some(resume);
546                } else {
547                    state.parser.reset();
548                    // after an incomplete value at the end there are no
549                    // more values
550                    state.ended = eof && err.kind() == ErrorKind::EndOfFile;
551                    state.failed = true;
552                }
553                Err(err)
554            }
555        }
556    }
557
558    fn peek(&mut self, input: &[u8], eof: bool) -> Result<Option<Progress>, Error> {
559        if !de::StreamDeserializer::supports_partial(self) {
560            return Ok(None);
561        }
562        Ok(Some(match self.skip_to_value(input, 0, eof)? {
563            Ok(pos) => Progress::Done { consumed: pos },
564            Err(progress) => progress,
565        }))
566    }
567}
568
569#[cfg(feature = "io")]
570impl DeserializerConfig {
571    /// Creates a reader of a stream of values (see
572    /// [`deser::io::Reader`](deser_core::io::Reader)).
573    ///
574    /// See [`StreamDeserializer`] for how the stream is split into values.
575    pub fn reader<R: Read>(&self, reader: R) -> deser_core::io::Reader<R, StreamDeserializer> {
576        deser_core::io::Reader::new(reader, StreamDeserializer::with_config(self.clone()))
577    }
578
579    /// Deserializes a value from a reader.
580    ///
581    /// See [`from_reader`].
582    pub fn from_reader<T: DeserializeOwned, R: Read>(&self, reader: R) -> Result<T, Error> {
583        deser_core::io::from_reader(reader, StreamDeserializer::with_config(self.clone()))
584    }
585}
586
587/// Deserializes a value from a reader.
588///
589/// The reader is read to the end.  Only whitespace may follow the value.
590/// The reader does not need to be buffered.  To read more than one value
591/// (for instance JSON Lines) use [`DeserializerConfig::reader`].
592///
593/// ```
594/// let value: Vec<u32> =
595///     deser_jsonc::from_reader(&b"[1, 2, 3]"[..]).unwrap();
596/// assert_eq!(value, [1, 2, 3]);
597/// ```
598#[cfg(feature = "io")]
599pub fn from_reader<T: DeserializeOwned, R: Read>(reader: R) -> Result<T, Error> {
600    DeserializerConfig::new().from_reader(reader)
601}