use crate::ast::{Token, Value};
use crate::error::{Error, Result};
#[cfg(feature = "serde")]
use serde::{Deserialize, Serialize};
use std::io::Read;
#[cfg(feature = "async")]
use tokio::io::{AsyncRead, AsyncReadExt};
pub trait JsonEventHandler: Send {
fn on_parse_start(&mut self) -> Result<()> {
Ok(())
}
fn on_object_start(&mut self) -> Result<()> {
Ok(())
}
fn on_object_end(&mut self) -> Result<()> {
Ok(())
}
fn on_array_start(&mut self) -> Result<()> {
Ok(())
}
fn on_array_end(&mut self) -> Result<()> {
Ok(())
}
fn on_key(&mut self, _key: &str) -> Result<()> {
Ok(())
}
fn on_value(&mut self, _value: &Value) -> Result<()> {
Ok(())
}
fn on_null(&mut self) -> Result<()> {
Ok(())
}
fn on_bool(&mut self, _value: bool) -> Result<()> {
Ok(())
}
fn on_number(&mut self, _value: &str) -> Result<()> {
Ok(())
}
fn on_string(&mut self, _value: &str) -> Result<()> {
Ok(())
}
fn on_parse_end(&mut self) -> Result<()> {
Ok(())
}
fn on_error(&mut self, error: &Error) -> Result<()> {
Err(error.clone())
}
}
#[derive(Debug, Clone)]
#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
pub struct ParserState {
pub position: usize,
pub context_stack: Vec<ParserContext>,
pub current_context: ParserContext,
pub is_complete: bool,
pub partial_buffer: String,
}
#[derive(Debug, Clone)]
#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
pub enum ParserContext {
Root,
Object {
expecting_key: bool,
current_key: Option<String>,
},
Array {
index: usize,
},
}
#[derive(Debug, Clone)]
pub struct EventParserConfig {
pub max_depth: usize,
pub chunk_size: usize,
pub json_paths: Vec<String>,
pub skip_large_values: bool,
pub large_value_threshold: usize,
}
impl Default for EventParserConfig {
fn default() -> Self {
EventParserConfig {
max_depth: 128,
chunk_size: 8192,
json_paths: Vec::new(),
skip_large_values: false,
large_value_threshold: 1024 * 1024, }
}
}
pub struct EventDrivenParser<H: JsonEventHandler> {
handler: H,
config: EventParserConfig,
state: ParserState,
path_matcher: Option<JsonPathMatcher>,
}
impl<H: JsonEventHandler> EventDrivenParser<H> {
pub fn new(handler: H) -> Self {
Self::with_config(handler, EventParserConfig::default())
}
pub fn with_config(handler: H, config: EventParserConfig) -> Self {
let path_matcher = if !config.json_paths.is_empty() {
Some(JsonPathMatcher::new(&config.json_paths))
} else {
None
};
EventDrivenParser {
handler,
config,
state: ParserState {
position: 0,
context_stack: Vec::new(),
current_context: ParserContext::Root,
is_complete: false,
partial_buffer: String::new(),
},
path_matcher,
}
}
pub fn parse<R: Read>(&mut self, reader: &mut R) -> Result<()> {
self.handler.on_parse_start()?;
let mut buffer = vec![0; self.config.chunk_size];
loop {
let bytes_read = reader
.read(&mut buffer)
.map_err(|e| Error::Custom(format!("IO error: {e}")))?;
if bytes_read == 0 {
break;
}
let chunk = std::str::from_utf8(&buffer[..bytes_read])
.map_err(|e| Error::Custom(format!("UTF-8 error: {e}")))?;
self.parse_chunk(chunk)?;
}
self.finish_parsing()?;
Ok(())
}
pub fn parse_chunk(&mut self, chunk: &str) -> Result<()> {
let combined_input;
let input = if self.state.partial_buffer.is_empty() {
chunk
} else {
combined_input = format!("{}{}", self.state.partial_buffer, chunk);
&combined_input
};
let (processed_bytes, remaining) = self.process_input(input)?;
self.state.position += processed_bytes;
self.state.partial_buffer = remaining;
Ok(())
}
fn process_input(&mut self, input: &str) -> Result<(usize, String)> {
let mut lexer = crate::lexer::Lexer::new(input);
let mut position = 0;
let mut in_object = false;
let mut expecting_key = false;
loop {
let (token, span) = lexer.next_token_with_span()?;
if !self.should_process_token(&token)? {
if token == Token::Eof {
break;
}
continue;
}
match &token {
Token::LeftBrace => {
self.handler.on_object_start()?;
self.push_context(ParserContext::Object {
expecting_key: true,
current_key: None,
});
in_object = true;
expecting_key = true;
}
Token::RightBrace => {
self.handler.on_object_end()?;
self.pop_context()?;
in_object = !self.state.context_stack.is_empty();
expecting_key = false;
}
Token::LeftBracket => {
self.handler.on_array_start()?;
self.push_context(ParserContext::Array { index: 0 });
}
Token::RightBracket => {
self.handler.on_array_end()?;
self.pop_context()?;
}
Token::String => {
let value = &input[span.start..span.end];
let unquoted = if value.starts_with('"') && value.ends_with('"') {
&value[1..value.len()-1]
} else {
value
};
if in_object && expecting_key {
self.handler.on_key(unquoted)?;
expecting_key = false;
} else {
self.handler.on_string(unquoted)?;
}
}
Token::Number => {
let value = &input[span.start..span.end];
self.handler.on_number(value)?;
}
Token::True => self.handler.on_bool(true)?,
Token::False => self.handler.on_bool(false)?,
Token::Null => self.handler.on_null()?,
Token::Comma => {
self.handle_comma()?;
if in_object {
expecting_key = true;
}
}
Token::Colon => {
self.handle_colon()?;
expecting_key = false;
}
Token::Eof => break,
_ => {}
}
position = span.end;
}
let remaining = if position < input.len() {
input[position..].to_string()
} else {
String::new()
};
Ok((position, remaining))
}
fn should_process_token(&self, _token: &Token) -> Result<bool> {
if let Some(ref matcher) = self.path_matcher {
Ok(matcher.matches(&self.get_current_path()))
} else {
Ok(true)
}
}
fn get_current_path(&self) -> String {
let mut path = String::from("$");
for context in &self.state.context_stack {
match context {
ParserContext::Object {
current_key: Some(key),
..
} => {
path.push('.');
path.push_str(key);
}
ParserContext::Array { index } => {
path.push_str(&format!("[{index}]"));
}
_ => {}
}
}
path
}
fn push_context(&mut self, context: ParserContext) {
if self.state.context_stack.len() >= self.config.max_depth {
return;
}
self.state
.context_stack
.push(self.state.current_context.clone());
self.state.current_context = context;
}
fn pop_context(&mut self) -> Result<()> {
if let Some(prev) = self.state.context_stack.pop() {
self.state.current_context = prev;
Ok(())
} else {
Err(Error::Custom("Unexpected closing bracket".to_string()))
}
}
#[allow(dead_code)]
fn update_object_context(&mut self, key: Option<String>) {
if let ParserContext::Object {
expecting_key,
current_key,
} = &mut self.state.current_context
{
*expecting_key = false;
*current_key = key;
}
}
fn handle_comma(&mut self) -> Result<()> {
match &mut self.state.current_context {
ParserContext::Object { expecting_key, .. } => {
*expecting_key = true;
}
ParserContext::Array { index } => {
*index += 1;
}
_ => {}
}
Ok(())
}
fn handle_colon(&mut self) -> Result<()> {
if let ParserContext::Object { expecting_key, .. } = &mut self.state.current_context {
*expecting_key = false;
}
Ok(())
}
pub fn finish_parsing(&mut self) -> Result<()> {
if !self.state.partial_buffer.is_empty() {
return Err(Error::Custom("Incomplete JSON at end of input".to_string()));
}
if !self.state.context_stack.is_empty() {
return Err(Error::Custom("Unclosed brackets".to_string()));
}
self.state.is_complete = true;
self.handler.on_parse_end()?;
Ok(())
}
pub fn save_state(&self) -> ParserState {
self.state.clone()
}
pub fn resume_from_state(mut self, state: ParserState) -> Self {
self.state = state;
self
}
}
struct JsonPathMatcher {
paths: Vec<String>,
}
impl JsonPathMatcher {
fn new(paths: &[String]) -> Self {
JsonPathMatcher {
paths: paths.to_vec(),
}
}
fn matches(&self, current_path: &str) -> bool {
self.paths.iter().any(|p| current_path.starts_with(p))
}
}
#[cfg(feature = "async")]
pub struct AsyncEventDrivenParser<H: JsonEventHandler> {
inner: EventDrivenParser<H>,
}
#[cfg(feature = "async")]
impl<H: JsonEventHandler> AsyncEventDrivenParser<H> {
pub fn new(handler: H) -> Self {
AsyncEventDrivenParser {
inner: EventDrivenParser::new(handler),
}
}
pub async fn parse<R: AsyncRead + Unpin>(&mut self, reader: &mut R) -> Result<()> {
self.inner.handler.on_parse_start()?;
let mut buffer = vec![0; self.inner.config.chunk_size];
loop {
let bytes_read = reader
.read(&mut buffer)
.await
.map_err(|e| Error::Custom(format!("Async IO error: {e}")))?;
if bytes_read == 0 {
break;
}
let chunk = std::str::from_utf8(&buffer[..bytes_read])
.map_err(|e| Error::Custom(format!("UTF-8 error: {e}")))?;
self.inner.parse_chunk(chunk)?;
}
self.inner.finish_parsing()?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
struct TestHandler {
events: Vec<String>,
}
impl JsonEventHandler for TestHandler {
fn on_object_start(&mut self) -> Result<()> {
self.events.push("object_start".to_string());
Ok(())
}
fn on_object_end(&mut self) -> Result<()> {
self.events.push("object_end".to_string());
Ok(())
}
fn on_key(&mut self, key: &str) -> Result<()> {
self.events.push(format!("key:{key}"));
Ok(())
}
fn on_string(&mut self, value: &str) -> Result<()> {
self.events.push(format!("string:{value}"));
Ok(())
}
}
#[test]
fn test_event_driven_parser() {
let handler = TestHandler { events: Vec::new() };
let mut parser = EventDrivenParser::new(handler);
let json = r#"{"name": "test", "value": "data"}"#;
let mut cursor = std::io::Cursor::new(json);
parser.parse(&mut cursor).unwrap();
let events = &parser.handler.events;
assert!(events.contains(&"object_start".to_string()));
assert!(events.contains(&"key:name".to_string()));
assert!(events.contains(&"string:test".to_string()));
assert!(events.contains(&"object_end".to_string()));
}
#[test]
fn test_resumable_parsing() {
let handler = TestHandler { events: Vec::new() };
let mut parser = EventDrivenParser::new(handler);
parser.parse_chunk(r#"{"name": "#).unwrap();
let state = parser.save_state();
assert!(!state.is_complete);
parser.parse_chunk(r#""test"}"#).unwrap();
parser.finish_parsing().unwrap();
assert!(parser.state.is_complete);
}
}