Skip to main content

large_json_array/
json_stream.rs

1use crate::error::JsonError;
2use serde_json::Value;
3use std::collections::VecDeque;
4use std::io::Read;
5
6pub struct JsonStream<R: Read> {
7    reader: R,
8    buffer: Vec<u8>,
9    temp: VecDeque<u8>, // 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: VecDeque::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        while let Some(c) = self.temp.pop_front() {
30            if self.brace_count > 0 || (c == b'{' && !self.in_string) {
31                self.object_buffer.push(char::from(c));
32            }
33
34            match c {
35                b'"' => {
36                    let mut backslashes = 0;
37                    while self.object_buffer.chars().rev().nth(backslashes) == Some('\\') {
38                        backslashes += 1;
39                    }
40                    if backslashes % 2 == 0 {
41                        self.in_string = !self.in_string;
42                    }
43                }
44                b'[' if !self.in_string && !self.inside_array => {
45                    self.inside_array = true;
46                    self.object_buffer.clear(); // drop the '['
47                    continue;
48                }
49                b'{' if !self.in_string => {
50                    self.brace_count += 1;
51                }
52                b'}' if !self.in_string => {
53                    self.brace_count -= 1;
54                    if self.brace_count == 0 {
55                        if self.object_buffer.trim().is_empty() {
56                            self.object_buffer.clear();
57                            continue;
58                        }
59
60                        let obj_str = self.object_buffer.clone();
61
62                        self.object_buffer.clear();
63
64                        while let Some(next_ch) = self.temp.front() {
65                            if next_ch.is_ascii_whitespace() || *next_ch == b',' {
66                                self.temp.pop_front();
67                            } else {
68                                break;
69                            }
70                        }
71                        return Some(serde_json::from_str(&obj_str).map_err(JsonError::from));
72                    }
73                }
74                _ => {}
75            }
76        }
77
78        None
79    }
80}
81
82impl<R: Read> Iterator for JsonStream<R> {
83    type Item = Result<Value, JsonError>;
84
85    fn next(&mut self) -> Option<Self::Item> {
86        if let Some(obj) = self.find_object_in_buffer() {
87            return Some(obj);
88        }
89        loop {
90            match self.reader.read(&mut self.buffer) {
91                Ok(0) => {
92                    return None;
93                }
94                Ok(n) => {
95                    self.temp.extend(&self.buffer[..n]);
96                    if let Some(obj) = self.find_object_in_buffer() {
97                        return Some(obj);
98                    }
99                }
100                Err(e) => {
101                    return Some(Err(JsonError::Io(e)));
102                }
103            }
104        }
105    }
106}