large_json_array/
json_stream.rs

1use crate::error::JsonError;
2use serde_json::Value;
3use std::io::Read;
4use std::str;
5
6pub struct JsonStream<R: Read> {
7    reader: R,
8    buffer: Vec<u8>,
9    temp: String, // leftover chars not yet processed
10    in_string: bool,
11    brace_count: u16,
12    inside_array: bool,
13    object_buffer: String, // characters of the current object
14}
15
16impl<R: Read> JsonStream<R> {
17    pub fn new(reader: R) -> Self {
18        JsonStream {
19            reader,
20            buffer: vec![0; 1024],
21            temp: String::new(),
22            in_string: false,
23            brace_count: 0,
24            inside_array: false,
25            object_buffer: String::new(),
26        }
27    }
28    pub fn find_object_in_buffer(&mut self) -> Option<Result<Value, JsonError>> {
29        let mut chars = self.temp.chars().peekable();
30        while let Some(c) = chars.next() {
31            if self.brace_count > 0 || (c == '{' && !self.in_string) {
32                self.object_buffer.push(c);
33            }
34
35            match c {
36                '"' => {
37                    let mut backslashes = 0;
38                    while self.object_buffer.chars().rev().nth(backslashes) == Some('\\') {
39                        backslashes += 1;
40                    }
41                    if backslashes % 2 == 0 {
42                        self.in_string = !self.in_string;
43                    }
44                }
45                '[' if !self.in_string && !self.inside_array => {
46                    self.inside_array = true;
47                    self.object_buffer.clear(); // drop the '['
48                    continue;
49                }
50                '{' if !self.in_string => {
51                    self.brace_count += 1;
52                }
53                '}' if !self.in_string => {
54                    self.brace_count -= 1;
55                    if self.brace_count == 0 {
56                        if self.object_buffer.trim().is_empty() {
57                            self.object_buffer.clear();
58                            continue;
59                        }
60
61                        let obj_str = self.object_buffer.clone();
62
63                        self.object_buffer.clear();
64
65                        while let Some(next_ch) = chars.peek() {
66                            if next_ch.is_whitespace() || *next_ch == ',' {
67                                chars.next();
68                            } else {
69                                break;
70                            }
71                        }
72
73                        // add remaining chars to buffer
74                        self.temp = chars.collect();
75
76                        return Some(serde_json::from_str(&obj_str).map_err(JsonError::from));
77                    }
78                }
79                _ => {}
80            }
81        }
82
83        // temp gets only unprocessed remainder
84        self.temp = chars.collect();
85        None
86    }
87}
88
89impl<R: Read> Iterator for JsonStream<R> {
90    type Item = Result<Value, JsonError>;
91
92    fn next(&mut self) -> Option<Self::Item> {
93        if let Some(obj) = self.find_object_in_buffer() {
94            return Some(obj);
95        }
96        while let Ok(n) = self.reader.read(&mut self.buffer) {
97            if n == 0 {
98                return None;
99            }
100            let chunk = str::from_utf8(&self.buffer[..n]).unwrap(); // assumes UTF-8 JSON
101            self.temp.push_str(chunk);
102            return self.find_object_in_buffer();
103        }
104        Some(Err(JsonError::Parser))
105    }
106}