mod buffered;
pub mod event_parser;
mod ndjson;
mod simple_lexer;
pub use buffered::{
parse_streaming, parse_streaming_with_config, BufferedStreamingConfig, BufferedStreamingParser,
StreamingEventIterator,
};
pub use event_parser::{
EventDrivenParser, EventParserConfig, JsonEventHandler, ParserContext as EventParserContext,
ParserState as EventParserState,
};
pub use ndjson::{NdJsonIterator, NdJsonParser, StreamingNdJsonParser};
pub use simple_lexer::SimpleStreamingLexer;
#[cfg(feature = "async")]
pub use event_parser::AsyncEventDrivenParser;
use crate::ast::{Token, Value};
use crate::error::{Error, Result};
use rustc_hash::FxHashMap;
#[derive(Debug, Clone, PartialEq)]
pub enum StreamingEvent {
StartObject,
EndObject,
StartArray,
EndArray,
ObjectKey(String),
Null,
Bool(bool),
Number(String),
String(String),
EndOfInput,
}
#[derive(Debug)]
pub struct StreamingParser {
lexer: SimpleStreamingLexer,
state_stack: Vec<ParserContext>,
current_state: ParserState,
options: crate::parser::ParserOptions,
event_queue: Vec<StreamingEvent>,
finished: bool,
current_token: Option<(Token, crate::error::Span)>,
input_buffer: String,
}
#[derive(Debug, Clone)]
enum ParserState {
ExpectingValue,
InObject { expecting_key: bool },
InArray {
#[allow(dead_code)]
first_element: bool,
},
BetweenValues,
ExpectingColon,
}
#[derive(Debug, Clone)]
enum ParserContext {
Object { expecting_key: bool },
Array { first_element: bool },
}
impl Default for StreamingParser {
fn default() -> Self {
Self::new()
}
}
impl StreamingParser {
pub fn new() -> Self {
Self::with_options(crate::parser::ParserOptions::default())
}
pub fn with_options(options: crate::parser::ParserOptions) -> Self {
let lexer = SimpleStreamingLexer::with_options(options.clone());
Self {
lexer,
state_stack: Vec::new(),
current_state: ParserState::ExpectingValue,
options,
event_queue: Vec::new(),
finished: false,
current_token: None,
input_buffer: String::new(),
}
}
pub fn feed(&mut self, chunk: &str) -> Result<()> {
if self.finished {
return Err(Error::Custom("Parser already finished".to_string()));
}
self.input_buffer.push_str(chunk);
self.lexer.feed_str(chunk)?;
self.process_tokens()?;
Ok(())
}
fn process_tokens(&mut self) -> Result<()> {
loop {
if self.current_token.is_none() {
self.current_token = self.lexer.next_token();
}
let Some((token, span)) = self.current_token else {
break;
};
if matches!(token, Token::SingleLineComment | Token::MultiLineComment) {
self.current_token = None;
continue;
}
if matches!(token, Token::Newline) {
self.process_newline(span)?;
self.current_token = None;
continue;
}
let consumed = match &self.current_state {
ParserState::ExpectingValue => self.process_value(token, span)?,
ParserState::InObject { expecting_key } => {
if *expecting_key {
self.process_object_key(token, span)?
} else {
self.process_value(token, span)?
}
}
ParserState::InArray { .. } => self.process_value(token, span)?,
ParserState::BetweenValues => self.process_between_values(token, span)?,
ParserState::ExpectingColon => {
if matches!(token, Token::Colon) {
self.current_state = ParserState::ExpectingValue;
true
} else {
return Err(Error::Expected {
expected: "colon".to_string(),
found: format!("{token:?}"),
position: span.start,
});
}
}
};
if consumed {
self.current_token = None;
} else {
break;
}
}
Ok(())
}
fn process_newline(&mut self, span: crate::error::Span) -> Result<()> {
if self.options.newline_as_comma && matches!(self.current_state, ParserState::BetweenValues) {
self.process_between_values(Token::Comma, span)?;
}
Ok(())
}
fn process_value(&mut self, token: Token, span: crate::error::Span) -> Result<bool> {
match token {
Token::LeftBrace => {
self.event_queue.push(StreamingEvent::StartObject);
self.state_stack.push(ParserContext::Object {
expecting_key: true,
});
self.current_state = ParserState::InObject {
expecting_key: true,
};
Ok(true)
}
Token::LeftBracket => {
self.event_queue.push(StreamingEvent::StartArray);
self.state_stack.push(ParserContext::Array {
first_element: true,
});
self.current_state = ParserState::InArray {
first_element: true,
};
Ok(true)
}
Token::String => {
let content = self.extract_string_content(span)?;
self.event_queue.push(StreamingEvent::String(content));
self.transition_after_value();
Ok(true)
}
Token::Number => {
let content = self.extract_token_content(span);
self.event_queue.push(StreamingEvent::Number(content));
self.transition_after_value();
Ok(true)
}
Token::True => {
self.event_queue.push(StreamingEvent::Bool(true));
self.transition_after_value();
Ok(true)
}
Token::False => {
self.event_queue.push(StreamingEvent::Bool(false));
self.transition_after_value();
Ok(true)
}
Token::Null => {
self.event_queue.push(StreamingEvent::Null);
self.transition_after_value();
Ok(true)
}
Token::RightBracket
if matches!(self.state_stack.last(), Some(ParserContext::Array { .. })) =>
{
self.event_queue.push(StreamingEvent::EndArray);
self.state_stack.pop();
self.transition_after_value();
Ok(true)
}
_ => Ok(false),
}
}
fn process_object_key(&mut self, token: Token, span: crate::error::Span) -> Result<bool> {
match token {
Token::String => {
let content = self.extract_string_content(span)?;
self.event_queue
.push(StreamingEvent::ObjectKey(content));
self.current_state = ParserState::ExpectingColon;
Ok(true)
}
Token::UnquotedString if self.options.allow_unquoted_keys => {
let content = self.extract_token_content(span);
self.event_queue
.push(StreamingEvent::ObjectKey(content));
self.current_state = ParserState::ExpectingColon;
Ok(true)
}
Token::RightBrace => {
self.event_queue.push(StreamingEvent::EndObject);
self.state_stack.pop();
self.transition_after_value();
Ok(true)
}
_ => Ok(false),
}
}
fn process_between_values(&mut self, token: Token, _span: crate::error::Span) -> Result<bool> {
match token {
Token::Comma => {
if let Some(context) = self.state_stack.last() {
match context {
ParserContext::Object { .. } => {
self.current_state = ParserState::InObject {
expecting_key: true,
};
}
ParserContext::Array { .. } => {
self.current_state = ParserState::InArray {
first_element: false,
};
}
}
}
Ok(true)
}
Token::RightBrace => {
if matches!(self.state_stack.last(), Some(ParserContext::Object { .. })) {
self.event_queue.push(StreamingEvent::EndObject);
self.state_stack.pop();
self.transition_after_value();
Ok(true)
} else {
Ok(false)
}
}
Token::RightBracket => {
if matches!(self.state_stack.last(), Some(ParserContext::Array { .. })) {
self.event_queue.push(StreamingEvent::EndArray);
self.state_stack.pop();
self.transition_after_value();
Ok(true)
} else {
Ok(false)
}
}
_ => Ok(false),
}
}
fn transition_after_value(&mut self) {
self.current_state = ParserState::BetweenValues;
}
pub fn next_event(&mut self) -> Result<Option<StreamingEvent>> {
self.process_tokens()?;
Ok(if self.event_queue.is_empty() {
None
} else {
Some(self.event_queue.remove(0))
})
}
pub fn finish(&mut self) -> Result<()> {
if self.finished {
return Ok(());
}
self.lexer.finish()?;
self.process_tokens()?;
if !self.state_stack.is_empty() {
return Err(Error::Custom("Unexpected end of input".to_string()));
}
self.event_queue.push(StreamingEvent::EndOfInput);
self.finished = true;
Ok(())
}
pub fn is_finished(&self) -> bool {
self.finished && self.event_queue.is_empty()
}
fn extract_token_content(&self, span: crate::error::Span) -> String {
if span.start < self.input_buffer.len() && span.end <= self.input_buffer.len() {
self.input_buffer[span.start..span.end].to_string()
} else {
String::new()
}
}
fn extract_string_content(&self, span: crate::error::Span) -> Result<String> {
let raw = self.extract_token_content(span);
let content = if (raw.starts_with('"') && raw.ends_with('"')) ||
(raw.starts_with('\'') && raw.ends_with('\'') && self.options.allow_single_quotes)
{
&raw[1..raw.len() - 1]
} else {
return Err(Error::Custom("Invalid string format".to_string()));
};
let mut result = String::new();
let mut chars = content.chars();
while let Some(ch) = chars.next() {
if ch == '\\' {
match chars.next() {
Some('"') => result.push('"'),
Some('\\') => result.push('\\'),
Some('/') => result.push('/'),
Some('b') => result.push('\u{0008}'),
Some('f') => result.push('\u{000C}'),
Some('n') => result.push('\n'),
Some('r') => result.push('\r'),
Some('t') => result.push('\t'),
Some('u') => {
let hex: String = chars.by_ref().take(4).collect();
if hex.len() != 4 {
return Err(Error::Custom("Invalid unicode escape".to_string()));
}
match u32::from_str_radix(&hex, 16) {
Ok(code) => {
if let Some(unicode_char) = char::from_u32(code) {
result.push(unicode_char);
} else {
return Err(Error::Custom(
"Invalid unicode code point".to_string(),
));
}
}
Err(_) => {
return Err(Error::Custom("Invalid unicode escape".to_string()))
}
}
}
Some(ch) => {
return Err(Error::Custom(format!("Invalid escape sequence: \\{ch}")))
}
None => return Err(Error::Custom("Incomplete escape sequence".to_string())),
}
} else {
result.push(ch);
}
}
Ok(result)
}
}
impl Iterator for StreamingParser {
type Item = Result<StreamingEvent>;
fn next(&mut self) -> Option<Self::Item> {
match self.next_event() {
Ok(Some(event)) => Some(Ok(event)),
Ok(None) => None,
Err(e) => Some(Err(e)),
}
}
}
pub struct StreamingValueBuilder {
stack: Vec<BuilderState>,
root: Option<Value>,
}
enum BuilderState {
Object(FxHashMap<String, Value>, Option<String>),
Array(Vec<Value>),
}
impl StreamingValueBuilder {
pub fn new() -> Self {
Self {
stack: Vec::new(),
root: None,
}
}
pub fn process_event(&mut self, event: StreamingEvent) -> Result<()> {
match event {
StreamingEvent::StartObject => {
self.stack
.push(BuilderState::Object(FxHashMap::default(), None));
}
StreamingEvent::StartArray => {
self.stack.push(BuilderState::Array(Vec::new()));
}
StreamingEvent::EndObject => {
if let Some(BuilderState::Object(map, _)) = self.stack.pop() {
self.add_value(Value::Object(map))?;
} else {
return Err(Error::Custom("Unexpected EndObject".to_string()));
}
}
StreamingEvent::EndArray => {
if let Some(BuilderState::Array(vec)) = self.stack.pop() {
self.add_value(Value::Array(vec))?;
} else {
return Err(Error::Custom("Unexpected EndArray".to_string()));
}
}
StreamingEvent::ObjectKey(key) => {
if let Some(BuilderState::Object(_, ref mut pending_key)) = self.stack.last_mut() {
*pending_key = Some(key);
} else {
return Err(Error::Custom("ObjectKey outside of object".to_string()));
}
}
StreamingEvent::Null => self.add_value(Value::Null)?,
StreamingEvent::Bool(b) => self.add_value(Value::Bool(b))?,
StreamingEvent::Number(n) => {
let value = if n.contains('.') || n.contains('e') || n.contains('E') {
Value::Number(crate::ast::Number::Float(
n.parse()
.map_err(|_| Error::Custom(format!("Invalid number: {n}")))?,
))
} else {
Value::Number(crate::ast::Number::Integer(
n.parse()
.map_err(|_| Error::Custom(format!("Invalid number: {n}")))?,
))
};
self.add_value(value)?;
}
StreamingEvent::String(s) => self.add_value(Value::String(s))?,
StreamingEvent::EndOfInput => {
if !self.stack.is_empty() {
return Err(Error::Custom("Unexpected end of input".to_string()));
}
}
}
Ok(())
}
fn add_value(&mut self, value: Value) -> Result<()> {
if self.stack.is_empty() {
if self.root.is_some() {
return Err(Error::Custom("Multiple root values".to_string()));
}
self.root = Some(value);
} else {
match self.stack.last_mut().unwrap() {
BuilderState::Object(map, pending_key) => {
if let Some(key) = pending_key.take() {
map.insert(key, value);
} else {
return Err(Error::Custom("Value without key in object".to_string()));
}
}
BuilderState::Array(vec) => {
vec.push(value);
}
}
}
Ok(())
}
pub fn finish(self) -> Result<Option<Value>> {
if !self.stack.is_empty() {
return Err(Error::Custom("Incomplete JSON structure".to_string()));
}
Ok(self.root)
}
}
impl Default for StreamingValueBuilder {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_streaming_parser_creation() {
let parser = StreamingParser::new();
assert!(!parser.is_finished());
}
#[test]
fn test_value_builder() {
let mut builder = StreamingValueBuilder::new();
builder.process_event(StreamingEvent::StartObject).unwrap();
builder
.process_event(StreamingEvent::ObjectKey("key".to_string()))
.unwrap();
builder
.process_event(StreamingEvent::String("value".to_string()))
.unwrap();
builder.process_event(StreamingEvent::EndObject).unwrap();
let value = builder.finish().unwrap().unwrap();
match value {
Value::Object(map) => {
assert_eq!(map.get("key").unwrap(), &Value::String("value".to_string()));
}
_ => panic!("Expected object"),
}
}
}