Skip to main content

vexy_json_core/streaming/
mod.rs

1// this_file: crates/core/src/streaming/mod.rs
2
3//! Streaming parser implementation for vexy_json.
4//!
5//! This module provides a streaming JSON parser that can process input
6//! incrementally, making it suitable for parsing large files or real-time
7//! data streams without loading the entire content into memory.
8
9mod 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/// Events emitted by the streaming parser
33#[derive(Debug, Clone, PartialEq)]
34pub enum StreamingEvent {
35    /// Start of a JSON object
36    StartObject,
37    /// End of a JSON object
38    EndObject,
39    /// Start of a JSON array
40    StartArray,
41    /// End of a JSON array
42    EndArray,
43    /// Object key (in objects, before the corresponding value)
44    ObjectKey(String),
45    /// Null value
46    Null,
47    /// Boolean value
48    Bool(bool),
49    /// Number value (as string to preserve precision)
50    Number(String),
51    /// String value
52    String(String),
53    /// End of input
54    EndOfInput,
55}
56
57/// State of the streaming parser
58#[derive(Debug)]
59pub struct StreamingParser {
60    /// Streaming lexer
61    lexer: SimpleStreamingLexer,
62    /// Parser state stack for nested structures
63    state_stack: Vec<ParserContext>,
64    /// Current parser state
65    current_state: ParserState,
66    /// Parser options
67    options: crate::parser::ParserOptions,
68    /// Event queue
69    event_queue: Vec<StreamingEvent>,
70    /// Whether parsing is complete
71    finished: bool,
72    /// Current token being processed
73    current_token: Option<(Token, crate::error::Span)>,
74    /// Input string for extracting token content
75    input_buffer: String,
76}
77
78/// Internal parser state
79#[derive(Debug, Clone)]
80enum ParserState {
81    /// Expecting a value (could be any JSON value)
82    ExpectingValue,
83    /// Inside an object, expecting key or closing brace
84    InObject { expecting_key: bool },
85    /// Inside an array, expecting value or closing bracket
86    InArray {
87        #[allow(dead_code)]
88        first_element: bool,
89    },
90    /// Between values (handling whitespace/commas)
91    BetweenValues,
92    /// Expecting a colon after object key
93    ExpectingColon,
94}
95
96/// Context for nested structures
97#[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    /// Create a new streaming parser with default options
111    pub fn new() -> Self {
112        Self::with_options(crate::parser::ParserOptions::default())
113    }
114
115    /// Create a new streaming parser with custom options
116    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    /// Feed a chunk of input to the parser
131    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        // Add to input buffer for token content extraction
137        self.input_buffer.push_str(chunk);
138
139        // Feed to lexer
140        self.lexer.feed_str(chunk)?;
141
142        // Process any available tokens
143        self.process_tokens()?;
144
145        Ok(())
146    }
147
148    /// Process tokens from the lexer
149    fn process_tokens(&mut self) -> Result<()> {
150        loop {
151            // Get next token if we don't have one
152            if self.current_token.is_none() {
153                self.current_token = self.lexer.next_token();
154            }
155
156            // If no token available, we're done for now
157            let Some((token, span)) = self.current_token else {
158                break;
159            };
160
161            // Skip comments
162            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            // Process token based on current state
174            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                // Token not consumed, stop processing
203                break;
204            }
205        }
206        Ok(())
207    }
208
209    /// Consume formatting newlines, using them as separators only after a value.
210    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    /// Process a value token
219    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                // Extract actual string content from input buffer
243                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                // Extract actual number content from input buffer
250                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    /// Process an object key
283    fn process_object_key(&mut self, token: Token, span: crate::error::Span) -> Result<bool> {
284        match token {
285            Token::String => {
286                // Extract actual string content from input buffer
287                let content = self.extract_string_content(span)?;
288                self.event_queue
289                    .push(StreamingEvent::ObjectKey(content));
290                // After key, expect colon then value
291                self.current_state = ParserState::ExpectingColon;
292                Ok(true)
293            }
294            Token::UnquotedString if self.options.allow_unquoted_keys => {
295                // Extract actual string content from input buffer
296                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                // Empty object or trailing comma
304                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    /// Process tokens between values
314    fn process_between_values(&mut self, token: Token, _span: crate::error::Span) -> Result<bool> {
315        match token {
316            Token::Comma => {
317                // Move to next value
318                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    /// Transition state after processing a value
359    fn transition_after_value(&mut self) {
360        self.current_state = ParserState::BetweenValues;
361    }
362
363    /// Get the next event from the parser
364    pub fn next_event(&mut self) -> Result<Option<StreamingEvent>> {
365        // Process more tokens if needed
366        self.process_tokens()?;
367
368        // Return queued event if available
369        Ok(if self.event_queue.is_empty() {
370            None
371        } else {
372            Some(self.event_queue.remove(0))
373        })
374    }
375
376    /// Signal end of input
377    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    /// Check if the parser has finished
395    pub fn is_finished(&self) -> bool {
396        self.finished && self.event_queue.is_empty()
397    }
398    
399    /// Extract token content from the input buffer
400    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            // Fallback for out-of-bounds
405            String::new()
406        }
407    }
408    
409    /// Extract string content from the input buffer, removing quotes and processing escapes
410    fn extract_string_content(&self, span: crate::error::Span) -> Result<String> {
411        let raw = self.extract_token_content(span);
412        
413        // Remove surrounding quotes
414        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        // Process escape sequences
423        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                        // Unicode escape sequence
439                        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
472/// Iterator interface for streaming events
473impl 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
485/// Builder pattern for constructing values from events
486pub 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    /// Create a new value builder
498    pub fn new() -> Self {
499        Self {
500            stack: Vec::new(),
501            root: None,
502        }
503    }
504
505    /// Process a streaming event
506    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                // Parse number string to Value::Number
540                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    /// Add a value to the current container
564    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    /// Get the final built value
588    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        // Build a simple object: {"key": "value"}
617        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}