1#[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;
18use crate::scan::skip_to_escape_single;
19
20fn skip_whitespace(input: &[u8], pos: usize, eof: bool) -> (usize, bool) {
26 let mut cursor = Cursor::new_partial(input, pos, eof);
27 let token = cursor.parse_whitespace().is_some();
28 (cursor.pos, token)
29}
30
31#[derive(Debug, Default)]
33struct StreamState {
34 done: bool,
36 parser: Parser,
38 skipping: Option<usize>,
41 failed: bool,
43 ended: bool,
45 pos: usize,
47 line: LineScan,
49 value: Option<Value>,
51}
52
53#[derive(Debug)]
55struct Value {
56 start: usize,
57 kind: ValueKind,
58 depth: usize,
59 in_string: bool,
60 single: bool,
62}
63
64#[derive(Debug)]
65enum ValueKind {
66 Scalar { parser: Parser, parsed: usize },
69 Structure,
71}
72
73fn frame_all(state: &mut StreamState, input: &[u8], eof: bool) -> Result<Frame, Error> {
74 if state.done {
75 return trailing_whitespace(input, 0, eof).map(|progress| match progress {
77 Progress::End => Frame::End,
78 Progress::NeedMore { consumed } => Frame::Incomplete { consumed },
80 Progress::Done { .. } => unreachable!(),
81 });
82 }
83 if !eof {
84 return Ok(Frame::Incomplete { consumed: 0 });
85 }
86 let (start, _) = skip_whitespace(input, 0, eof);
88 Ok(if start < input.len() {
89 state.done = true;
90 Frame::Value {
91 start,
92 end: input.len(),
93 consumed: input.len(),
94 }
95 } else {
96 Frame::End
97 })
98}
99
100fn trailing_whitespace(input: &[u8], offset: usize, eof: bool) -> Result<Progress, Error> {
104 match skip_whitespace(input, 0, eof) {
105 (pos, _) if pos == input.len() && eof => Ok(Progress::End),
106 (pos, false) if !eof => Ok(Progress::NeedMore { consumed: pos }),
108 (pos, _) => Err(Error::with_offset(
109 ErrorKind::Syntax,
110 "garbage after input",
111 offset + pos,
112 )),
113 }
114}
115
116fn frame_line(state: &mut StreamState, input: &[u8], eof: bool) -> Frame {
117 if state.pos == 0 {
120 match skip_whitespace(input, 0, eof) {
121 (0, true) => {}
122 (start, true) => return Frame::Incomplete { consumed: start },
123 (_, false) if input.is_empty() && eof => return Frame::End,
124 (start, false) if eof && start < input.len() => {
126 return Frame::Value {
127 start,
128 end: input.len(),
129 consumed: input.len(),
130 };
131 }
132 (start, false) => return Frame::Incomplete { consumed: start },
134 }
135 }
136 let end = state.line.find_end(input, state.pos);
138 let end = match end {
139 Some(end) => end,
140 None if eof => input.len(),
141 None => {
142 state.pos = input.len();
143 return Frame::Incomplete { consumed: 0 };
144 }
145 };
146 state.pos = 0;
147 state.line = LineScan::default();
148 let consumed = (end + 1).min(input.len());
149 match skip_whitespace(&input[..end], 0, true) {
150 (start, _) if start < end => Frame::Value {
151 start,
152 end,
153 consumed,
154 },
155 _ if consumed == 0 => Frame::End,
156 _ => Frame::Incomplete { consumed },
158 }
159}
160
161fn frame_value(state: &mut StreamState, input: &[u8], eof: bool) -> Frame {
162 let value = match state.value {
163 Some(ref mut value) => value,
164 None => {
165 let start = match skip_whitespace(input, 0, eof) {
166 (start, true) => start,
167 (_, false) if input.is_empty() && eof => return Frame::End,
168 (consumed, false) => return Frame::Incomplete { consumed },
169 };
170 let (kind, depth, in_string) = match input[start] {
171 b'"' => (ValueKind::Structure, 0, true),
172 b'\'' => (ValueKind::Structure, 0, true),
173 b'{' | b'[' => (ValueKind::Structure, 1, false),
174 b'}' | b']' | b',' | b':' => {
177 return Frame::Value {
178 start,
179 end: start + 1,
180 consumed: start + 1,
181 };
182 }
183 _ => (
184 ValueKind::Scalar {
185 parser: Parser::default(),
186 parsed: 0,
187 },
188 0,
189 false,
190 ),
191 };
192 state.pos = start + 1;
193 state.value.insert(Value {
194 start,
195 kind,
196 depth,
197 in_string,
198 single: input[start] == b'\'',
199 })
200 }
201 };
202
203 let end = match value.kind {
204 ValueKind::Scalar {
205 ref mut parser,
206 ref mut parsed,
207 } => {
208 let pos = value.start + *parsed;
209 let options = Options {
210 validate_utf8: true,
211 exact_numbers: false,
212 };
213 let mut discard = Discard(State::new());
214 match parser.parse(&input[pos..], 0, eof, 0, options, &mut discard) {
215 Ok(ParseProgress::Done(end)) => Some(pos + end),
216 Ok(ParseProgress::NeedMore(consumed)) => {
217 *parsed += consumed;
218 state.pos = input.len();
219 None
220 }
221 Err(err) => Some(error_end(&err, input, pos)),
224 }
225 }
226 ValueKind::Structure => scan_structure(input, &mut state.pos, value),
227 };
228
229 let start = value.start;
230 match end {
231 Some(end) => {
232 state.value = None;
233 state.pos = 0;
234 Frame::Value {
235 start,
236 end,
237 consumed: end,
238 }
239 }
240 None if eof => {
242 state.value = None;
243 state.pos = 0;
244 Frame::Value {
245 start,
246 end: input.len(),
247 consumed: input.len(),
248 }
249 }
250 None => {
251 value.start = 0;
253 state.pos -= start;
254 Frame::Incomplete { consumed: start }
255 }
256 }
257}
258
259fn error_end(err: &Error, input: &[u8], pos: usize) -> usize {
267 match err.offset() {
268 Some(offset) if err.kind() != ErrorKind::EndOfFile => (pos + offset + 1).min(input.len()),
269 _ => input.len(),
270 }
271}
272
273fn scan_structure(input: &[u8], pos: &mut usize, value: &mut Value) -> Option<usize> {
277 let mut index = *pos;
278 while index < input.len() {
279 if value.in_string {
280 index = if value.single {
281 skip_to_escape_single(input, index)
282 } else {
283 skip_to_escape(input, index)
284 };
285 let byte = input.get(index).copied();
286 let byte = if value.single && byte == Some(b'\'') {
288 Some(b'"')
289 } else {
290 byte
291 };
292 match byte {
293 Some(b'"') => {
294 value.in_string = false;
295 index += 1;
296 if value.depth == 0 {
297 return Some(index);
298 }
299 }
300 Some(b'\\') => {
301 if index + 1 >= input.len() {
303 break;
304 }
305 index += 2;
306 }
307 Some(_) => index += 1,
309 None => break,
310 }
311 } else {
312 match input[index] {
313 b'"' => {
314 value.in_string = true;
315 value.single = false;
316 }
317 b'\'' => {
318 value.in_string = true;
319 value.single = true;
320 }
321 b'/' => {
324 match skip_whitespace(input, index, false) {
325 (next, true) if next == index => index += 1,
327 (next, token) => {
328 index = next;
329 if !token && next < input.len() {
330 break;
331 }
332 }
333 }
334 continue;
335 }
336 b'{' | b'[' => value.depth += 1,
337 b'}' | b']' => {
338 value.depth -= 1;
339 if value.depth == 0 {
340 return Some(index + 1);
341 }
342 }
343 _ => {}
344 }
345 index += 1;
346 }
347 }
348 *pos = index;
349 None
350}
351
352#[derive(Debug)]
393pub struct StreamDeserializer {
394 config: DeserializerConfig,
395 state: StreamState,
396}
397
398impl Default for StreamDeserializer {
399 fn default() -> StreamDeserializer {
400 StreamDeserializer::new()
401 }
402}
403
404impl StreamDeserializer {
405 pub fn new() -> StreamDeserializer {
407 StreamDeserializer::with_config(DeserializerConfig::new())
408 }
409
410 pub fn with_config(config: DeserializerConfig) -> StreamDeserializer {
412 StreamDeserializer {
413 config,
414 state: StreamState::default(),
415 }
416 }
417
418 pub fn config(&self) -> &DeserializerConfig {
420 &self.config
421 }
422
423 fn skip_to_value(
429 &mut self,
430 input: &[u8],
431 offset: usize,
432 eof: bool,
433 ) -> Result<Result<usize, Progress>, Error> {
434 let options = self.options();
435 let state = &mut self.state;
436 if state.ended {
437 return Ok(Err(Progress::End));
438 }
439 if state.failed {
440 return Err(Error::new(
441 ErrorKind::InvalidState,
442 "cannot continue after an error",
443 ));
444 }
445
446 let mut pos = 0;
448 if let Some(skip) = state.skipping {
449 let mut discard = Discard(State::new());
450 match state
451 .parser
452 .parse(input, skip, eof, offset, options, &mut discard)
453 {
454 Ok(ParseProgress::Done(end)) => {
455 state.skipping = None;
456 state.done = self.config.trailing_mode() == Trailing::Strict;
457 pos = end;
458 }
459 Ok(ParseProgress::NeedMore(consumed)) => {
460 state.skipping = Some(0);
461 return Ok(Err(Progress::NeedMore { consumed }));
462 }
463 Err(err) => {
464 state.parser.reset();
465 state.skipping = None;
466 state.failed = true;
467 return Err(err);
468 }
469 }
470 }
471
472 if !state.parser.is_idle() {
473 return Ok(Ok(pos));
474 }
475 if state.done {
476 return match trailing_whitespace(&input[pos..], offset + pos, eof)? {
477 Progress::NeedMore { consumed } => Ok(Err(Progress::NeedMore {
478 consumed: pos + consumed,
479 })),
480 progress => Ok(Err(progress)),
481 };
482 }
483 let token;
486 (pos, token) = skip_whitespace(input, pos, eof);
487 if !token && (pos == input.len() || !eof) {
488 return Ok(Err(if eof {
489 Progress::End
490 } else {
491 Progress::NeedMore { consumed: pos }
492 }));
493 }
494 Ok(Ok(pos))
495 }
496
497 fn options(&self) -> Options {
499 Options {
500 validate_utf8: true,
501 exact_numbers: self.config.exact_numbers_enabled(),
502 }
503 }
504}
505
506impl de::StreamDeserializer for StreamDeserializer {
507 fn context(&self) -> deser_core::Context {
508 self.config.context().clone()
509 }
510
511 fn frame(&mut self, input: &[u8], eof: bool) -> Result<Frame, Error> {
512 let state = &mut self.state;
513 match self.config.trailing_mode() {
514 Trailing::Strict => frame_all(state, input, eof),
515 Trailing::Newline => Ok(frame_line(state, input, eof)),
516 Trailing::Stop => Ok(frame_value(state, input, eof)),
517 }
518 }
519
520 fn drive_frame<'de>(
521 &mut self,
522 frame: &'de [u8],
523 driver: &mut DeserializeDriver<'_, 'de>,
524 ) -> Result<(), Error> {
525 let mut de = Deserializer::from_frame(frame, &self.config);
526 de.drive(driver)
527 }
528
529 fn is_text(&self) -> bool {
530 true
531 }
532
533 fn supports_partial(&self) -> bool {
536 self.config.trailing_mode() != Trailing::Newline
537 }
538
539 fn drive_partial(
540 &mut self,
541 input: &[u8],
542 offset: usize,
543 eof: bool,
544 driver: &mut DeserializeDriver<'_, '_>,
545 ) -> Result<Progress, Error> {
546 let pos = match self.skip_to_value(input, offset, eof)? {
547 Ok(pos) => pos,
548 Err(progress) => return Ok(progress),
549 };
550 let options = self.options();
551 let state = &mut self.state;
552 match state
553 .parser
554 .parse(input, pos, eof, offset, options, &mut Copying(driver))
555 {
556 Ok(ParseProgress::Done(end)) => {
557 state.done = self.config.trailing_mode() == Trailing::Strict;
558 Ok(Progress::Done { consumed: end })
559 }
560 Ok(ParseProgress::NeedMore(consumed)) => Ok(Progress::NeedMore { consumed }),
561 Err(err) => {
562 if let Some(resume) = state.parser.recoverable() {
563 state.skipping = Some(resume);
566 } else {
567 state.parser.reset();
568 state.ended = eof && err.kind() == ErrorKind::EndOfFile;
571 state.failed = true;
572 }
573 Err(err)
574 }
575 }
576 }
577
578 fn peek(&mut self, input: &[u8], eof: bool) -> Result<Option<Progress>, Error> {
579 if !de::StreamDeserializer::supports_partial(self) {
580 return Ok(None);
581 }
582 Ok(Some(match self.skip_to_value(input, 0, eof)? {
583 Ok(pos) => Progress::Done { consumed: pos },
584 Err(progress) => progress,
585 }))
586 }
587}
588
589#[cfg(feature = "io")]
590impl DeserializerConfig {
591 pub fn reader<R: Read>(&self, reader: R) -> deser_core::io::Reader<R, StreamDeserializer> {
596 deser_core::io::Reader::new(reader, StreamDeserializer::with_config(self.clone()))
597 }
598
599 pub fn from_reader<T: DeserializeOwned, R: Read>(&self, reader: R) -> Result<T, Error> {
603 deser_core::io::from_reader(reader, StreamDeserializer::with_config(self.clone()))
604 }
605}
606
607#[cfg(feature = "io")]
619pub fn from_reader<T: DeserializeOwned, R: Read>(reader: R) -> Result<T, Error> {
620 DeserializerConfig::new().from_reader(reader)
621}