vexy_json_core/streaming/
mod.rs1mod buffered;
10pub mod event_parser;
11mod ndjson;
12mod simple_lexer;
13
14pub use buffered::{
15 parse_streaming, parse_streaming_with_config, BufferedStreamingConfig, BufferedStreamingParser,
16 StreamingEventIterator,
17};
18pub use event_parser::{
19 EventDrivenParser, EventParserConfig, JsonEventHandler, ParserContext as EventParserContext,
20 ParserState as EventParserState,
21};
22pub use ndjson::{NdJsonIterator, NdJsonParser, StreamingNdJsonParser};
23pub use simple_lexer::SimpleStreamingLexer;
24
25#[cfg(feature = "async")]
26pub use event_parser::AsyncEventDrivenParser;
27
28use crate::ast::{Token, Value};
29use crate::error::{Error, Result};
30use rustc_hash::FxHashMap;
31
32#[derive(Debug, Clone, PartialEq)]
34pub enum StreamingEvent {
35 StartObject,
37 EndObject,
39 StartArray,
41 EndArray,
43 ObjectKey(String),
45 Null,
47 Bool(bool),
49 Number(String),
51 String(String),
53 EndOfInput,
55}
56
57#[derive(Debug)]
59pub struct StreamingParser {
60 lexer: SimpleStreamingLexer,
62 state_stack: Vec<ParserContext>,
64 current_state: ParserState,
66 options: crate::parser::ParserOptions,
68 event_queue: Vec<StreamingEvent>,
70 finished: bool,
72 current_token: Option<(Token, crate::error::Span)>,
74 input_buffer: String,
76}
77
78#[derive(Debug, Clone)]
80enum ParserState {
81 ExpectingValue,
83 InObject { expecting_key: bool },
85 InArray {
87 #[allow(dead_code)]
88 first_element: bool,
89 },
90 BetweenValues,
92 ExpectingColon,
94}
95
96#[derive(Debug, Clone)]
98enum ParserContext {
99 Object { expecting_key: bool },
100 Array { first_element: bool },
101}
102
103impl Default for StreamingParser {
104 fn default() -> Self {
105 Self::new()
106 }
107}
108
109impl StreamingParser {
110 pub fn new() -> Self {
112 Self::with_options(crate::parser::ParserOptions::default())
113 }
114
115 pub fn with_options(options: crate::parser::ParserOptions) -> Self {
117 let lexer = SimpleStreamingLexer::with_options(options.clone());
118 Self {
119 lexer,
120 state_stack: Vec::new(),
121 current_state: ParserState::ExpectingValue,
122 options,
123 event_queue: Vec::new(),
124 finished: false,
125 current_token: None,
126 input_buffer: String::new(),
127 }
128 }
129
130 pub fn feed(&mut self, chunk: &str) -> Result<()> {
132 if self.finished {
133 return Err(Error::Custom("Parser already finished".to_string()));
134 }
135
136 self.input_buffer.push_str(chunk);
138
139 self.lexer.feed_str(chunk)?;
141
142 self.process_tokens()?;
144
145 Ok(())
146 }
147
148 fn process_tokens(&mut self) -> Result<()> {
150 loop {
151 if self.current_token.is_none() {
153 self.current_token = self.lexer.next_token();
154 }
155
156 let Some((token, span)) = self.current_token else {
158 break;
159 };
160
161 if matches!(token, Token::SingleLineComment | Token::MultiLineComment) {
163 self.current_token = None;
164 continue;
165 }
166
167 if matches!(token, Token::Newline) {
168 self.process_newline(span)?;
169 self.current_token = None;
170 continue;
171 }
172
173 let consumed = match &self.current_state {
175 ParserState::ExpectingValue => self.process_value(token, span)?,
176 ParserState::InObject { expecting_key } => {
177 if *expecting_key {
178 self.process_object_key(token, span)?
179 } else {
180 self.process_value(token, span)?
181 }
182 }
183 ParserState::InArray { .. } => self.process_value(token, span)?,
184 ParserState::BetweenValues => self.process_between_values(token, span)?,
185 ParserState::ExpectingColon => {
186 if matches!(token, Token::Colon) {
187 self.current_state = ParserState::ExpectingValue;
188 true
189 } else {
190 return Err(Error::Expected {
191 expected: "colon".to_string(),
192 found: format!("{token:?}"),
193 position: span.start,
194 });
195 }
196 }
197 };
198
199 if consumed {
200 self.current_token = None;
201 } else {
202 break;
204 }
205 }
206 Ok(())
207 }
208
209 fn process_newline(&mut self, span: crate::error::Span) -> Result<()> {
211 if self.options.newline_as_comma && matches!(self.current_state, ParserState::BetweenValues) {
212 self.process_between_values(Token::Comma, span)?;
213 }
214
215 Ok(())
216 }
217
218 fn process_value(&mut self, token: Token, span: crate::error::Span) -> Result<bool> {
220 match token {
221 Token::LeftBrace => {
222 self.event_queue.push(StreamingEvent::StartObject);
223 self.state_stack.push(ParserContext::Object {
224 expecting_key: true,
225 });
226 self.current_state = ParserState::InObject {
227 expecting_key: true,
228 };
229 Ok(true)
230 }
231 Token::LeftBracket => {
232 self.event_queue.push(StreamingEvent::StartArray);
233 self.state_stack.push(ParserContext::Array {
234 first_element: true,
235 });
236 self.current_state = ParserState::InArray {
237 first_element: true,
238 };
239 Ok(true)
240 }
241 Token::String => {
242 let content = self.extract_string_content(span)?;
244 self.event_queue.push(StreamingEvent::String(content));
245 self.transition_after_value();
246 Ok(true)
247 }
248 Token::Number => {
249 let content = self.extract_token_content(span);
251 self.event_queue.push(StreamingEvent::Number(content));
252 self.transition_after_value();
253 Ok(true)
254 }
255 Token::True => {
256 self.event_queue.push(StreamingEvent::Bool(true));
257 self.transition_after_value();
258 Ok(true)
259 }
260 Token::False => {
261 self.event_queue.push(StreamingEvent::Bool(false));
262 self.transition_after_value();
263 Ok(true)
264 }
265 Token::Null => {
266 self.event_queue.push(StreamingEvent::Null);
267 self.transition_after_value();
268 Ok(true)
269 }
270 Token::RightBracket
271 if matches!(self.state_stack.last(), Some(ParserContext::Array { .. })) =>
272 {
273 self.event_queue.push(StreamingEvent::EndArray);
274 self.state_stack.pop();
275 self.transition_after_value();
276 Ok(true)
277 }
278 _ => Ok(false),
279 }
280 }
281
282 fn process_object_key(&mut self, token: Token, span: crate::error::Span) -> Result<bool> {
284 match token {
285 Token::String => {
286 let content = self.extract_string_content(span)?;
288 self.event_queue
289 .push(StreamingEvent::ObjectKey(content));
290 self.current_state = ParserState::ExpectingColon;
292 Ok(true)
293 }
294 Token::UnquotedString if self.options.allow_unquoted_keys => {
295 let content = self.extract_token_content(span);
297 self.event_queue
298 .push(StreamingEvent::ObjectKey(content));
299 self.current_state = ParserState::ExpectingColon;
300 Ok(true)
301 }
302 Token::RightBrace => {
303 self.event_queue.push(StreamingEvent::EndObject);
305 self.state_stack.pop();
306 self.transition_after_value();
307 Ok(true)
308 }
309 _ => Ok(false),
310 }
311 }
312
313 fn process_between_values(&mut self, token: Token, _span: crate::error::Span) -> Result<bool> {
315 match token {
316 Token::Comma => {
317 if let Some(context) = self.state_stack.last() {
319 match context {
320 ParserContext::Object { .. } => {
321 self.current_state = ParserState::InObject {
322 expecting_key: true,
323 };
324 }
325 ParserContext::Array { .. } => {
326 self.current_state = ParserState::InArray {
327 first_element: false,
328 };
329 }
330 }
331 }
332 Ok(true)
333 }
334 Token::RightBrace => {
335 if matches!(self.state_stack.last(), Some(ParserContext::Object { .. })) {
336 self.event_queue.push(StreamingEvent::EndObject);
337 self.state_stack.pop();
338 self.transition_after_value();
339 Ok(true)
340 } else {
341 Ok(false)
342 }
343 }
344 Token::RightBracket => {
345 if matches!(self.state_stack.last(), Some(ParserContext::Array { .. })) {
346 self.event_queue.push(StreamingEvent::EndArray);
347 self.state_stack.pop();
348 self.transition_after_value();
349 Ok(true)
350 } else {
351 Ok(false)
352 }
353 }
354 _ => Ok(false),
355 }
356 }
357
358 fn transition_after_value(&mut self) {
360 self.current_state = ParserState::BetweenValues;
361 }
362
363 pub fn next_event(&mut self) -> Result<Option<StreamingEvent>> {
365 self.process_tokens()?;
367
368 Ok(if self.event_queue.is_empty() {
370 None
371 } else {
372 Some(self.event_queue.remove(0))
373 })
374 }
375
376 pub fn finish(&mut self) -> Result<()> {
378 if self.finished {
379 return Ok(());
380 }
381
382 self.lexer.finish()?;
383 self.process_tokens()?;
384
385 if !self.state_stack.is_empty() {
386 return Err(Error::Custom("Unexpected end of input".to_string()));
387 }
388
389 self.event_queue.push(StreamingEvent::EndOfInput);
390 self.finished = true;
391 Ok(())
392 }
393
394 pub fn is_finished(&self) -> bool {
396 self.finished && self.event_queue.is_empty()
397 }
398
399 fn extract_token_content(&self, span: crate::error::Span) -> String {
401 if span.start < self.input_buffer.len() && span.end <= self.input_buffer.len() {
402 self.input_buffer[span.start..span.end].to_string()
403 } else {
404 String::new()
406 }
407 }
408
409 fn extract_string_content(&self, span: crate::error::Span) -> Result<String> {
411 let raw = self.extract_token_content(span);
412
413 let content = if (raw.starts_with('"') && raw.ends_with('"')) ||
415 (raw.starts_with('\'') && raw.ends_with('\'') && self.options.allow_single_quotes)
416 {
417 &raw[1..raw.len() - 1]
418 } else {
419 return Err(Error::Custom("Invalid string format".to_string()));
420 };
421
422 let mut result = String::new();
424 let mut chars = content.chars();
425
426 while let Some(ch) = chars.next() {
427 if ch == '\\' {
428 match chars.next() {
429 Some('"') => result.push('"'),
430 Some('\\') => result.push('\\'),
431 Some('/') => result.push('/'),
432 Some('b') => result.push('\u{0008}'),
433 Some('f') => result.push('\u{000C}'),
434 Some('n') => result.push('\n'),
435 Some('r') => result.push('\r'),
436 Some('t') => result.push('\t'),
437 Some('u') => {
438 let hex: String = chars.by_ref().take(4).collect();
440 if hex.len() != 4 {
441 return Err(Error::Custom("Invalid unicode escape".to_string()));
442 }
443 match u32::from_str_radix(&hex, 16) {
444 Ok(code) => {
445 if let Some(unicode_char) = char::from_u32(code) {
446 result.push(unicode_char);
447 } else {
448 return Err(Error::Custom(
449 "Invalid unicode code point".to_string(),
450 ));
451 }
452 }
453 Err(_) => {
454 return Err(Error::Custom("Invalid unicode escape".to_string()))
455 }
456 }
457 }
458 Some(ch) => {
459 return Err(Error::Custom(format!("Invalid escape sequence: \\{ch}")))
460 }
461 None => return Err(Error::Custom("Incomplete escape sequence".to_string())),
462 }
463 } else {
464 result.push(ch);
465 }
466 }
467
468 Ok(result)
469 }
470}
471
472impl Iterator for StreamingParser {
474 type Item = Result<StreamingEvent>;
475
476 fn next(&mut self) -> Option<Self::Item> {
477 match self.next_event() {
478 Ok(Some(event)) => Some(Ok(event)),
479 Ok(None) => None,
480 Err(e) => Some(Err(e)),
481 }
482 }
483}
484
485pub struct StreamingValueBuilder {
487 stack: Vec<BuilderState>,
488 root: Option<Value>,
489}
490
491enum BuilderState {
492 Object(FxHashMap<String, Value>, Option<String>),
493 Array(Vec<Value>),
494}
495
496impl StreamingValueBuilder {
497 pub fn new() -> Self {
499 Self {
500 stack: Vec::new(),
501 root: None,
502 }
503 }
504
505 pub fn process_event(&mut self, event: StreamingEvent) -> Result<()> {
507 match event {
508 StreamingEvent::StartObject => {
509 self.stack
510 .push(BuilderState::Object(FxHashMap::default(), None));
511 }
512 StreamingEvent::StartArray => {
513 self.stack.push(BuilderState::Array(Vec::new()));
514 }
515 StreamingEvent::EndObject => {
516 if let Some(BuilderState::Object(map, _)) = self.stack.pop() {
517 self.add_value(Value::Object(map))?;
518 } else {
519 return Err(Error::Custom("Unexpected EndObject".to_string()));
520 }
521 }
522 StreamingEvent::EndArray => {
523 if let Some(BuilderState::Array(vec)) = self.stack.pop() {
524 self.add_value(Value::Array(vec))?;
525 } else {
526 return Err(Error::Custom("Unexpected EndArray".to_string()));
527 }
528 }
529 StreamingEvent::ObjectKey(key) => {
530 if let Some(BuilderState::Object(_, ref mut pending_key)) = self.stack.last_mut() {
531 *pending_key = Some(key);
532 } else {
533 return Err(Error::Custom("ObjectKey outside of object".to_string()));
534 }
535 }
536 StreamingEvent::Null => self.add_value(Value::Null)?,
537 StreamingEvent::Bool(b) => self.add_value(Value::Bool(b))?,
538 StreamingEvent::Number(n) => {
539 let value = if n.contains('.') || n.contains('e') || n.contains('E') {
541 Value::Number(crate::ast::Number::Float(
542 n.parse()
543 .map_err(|_| Error::Custom(format!("Invalid number: {n}")))?,
544 ))
545 } else {
546 Value::Number(crate::ast::Number::Integer(
547 n.parse()
548 .map_err(|_| Error::Custom(format!("Invalid number: {n}")))?,
549 ))
550 };
551 self.add_value(value)?;
552 }
553 StreamingEvent::String(s) => self.add_value(Value::String(s))?,
554 StreamingEvent::EndOfInput => {
555 if !self.stack.is_empty() {
556 return Err(Error::Custom("Unexpected end of input".to_string()));
557 }
558 }
559 }
560 Ok(())
561 }
562
563 fn add_value(&mut self, value: Value) -> Result<()> {
565 if self.stack.is_empty() {
566 if self.root.is_some() {
567 return Err(Error::Custom("Multiple root values".to_string()));
568 }
569 self.root = Some(value);
570 } else {
571 match self.stack.last_mut().unwrap() {
572 BuilderState::Object(map, pending_key) => {
573 if let Some(key) = pending_key.take() {
574 map.insert(key, value);
575 } else {
576 return Err(Error::Custom("Value without key in object".to_string()));
577 }
578 }
579 BuilderState::Array(vec) => {
580 vec.push(value);
581 }
582 }
583 }
584 Ok(())
585 }
586
587 pub fn finish(self) -> Result<Option<Value>> {
589 if !self.stack.is_empty() {
590 return Err(Error::Custom("Incomplete JSON structure".to_string()));
591 }
592 Ok(self.root)
593 }
594}
595
596impl Default for StreamingValueBuilder {
597 fn default() -> Self {
598 Self::new()
599 }
600}
601
602#[cfg(test)]
603mod tests {
604 use super::*;
605
606 #[test]
607 fn test_streaming_parser_creation() {
608 let parser = StreamingParser::new();
609 assert!(!parser.is_finished());
610 }
611
612 #[test]
613 fn test_value_builder() {
614 let mut builder = StreamingValueBuilder::new();
615
616 builder.process_event(StreamingEvent::StartObject).unwrap();
618 builder
619 .process_event(StreamingEvent::ObjectKey("key".to_string()))
620 .unwrap();
621 builder
622 .process_event(StreamingEvent::String("value".to_string()))
623 .unwrap();
624 builder.process_event(StreamingEvent::EndObject).unwrap();
625
626 let value = builder.finish().unwrap().unwrap();
627 match value {
628 Value::Object(map) => {
629 assert_eq!(map.get("key").unwrap(), &Value::String("value".to_string()));
630 }
631 _ => panic!("Expected object"),
632 }
633 }
634}