use serde::de::DeserializeOwned;
use schemars::JsonSchema;
use super::json_utils::{find_json_structures, deserialize_stream_map, ParsedOrUnknown, JsonStreamParser};
use super::{StreamItem, TextContent};
use tracing::debug;
#[derive(Debug)]
pub struct JsonStreamProcessor<T> {
parser: JsonStreamParser,
text_buf: String,
last_offset: usize,
_phantom: std::marker::PhantomData<T>,
}
impl<T> JsonStreamProcessor<T>
where
T: DeserializeOwned + JsonSchema + Send + 'static,
{
pub fn new() -> Self {
Self {
parser: JsonStreamParser::new(),
text_buf: String::new(),
last_offset: 0,
_phantom: std::marker::PhantomData,
}
}
pub fn process_chunk(&mut self, chunk: &str) -> Vec<StreamItem<T>> {
let mut items = Vec::new();
self.text_buf.push_str(chunk);
for node in self.parser.feed(chunk) {
if node.start > self.last_offset && node.start <= self.text_buf.len() {
let text_slice = &self.text_buf[self.last_offset..node.start];
if !text_slice.trim().is_empty() {
items.push(StreamItem::Text(TextContent { text: text_slice.to_string() }));
}
}
let end = node.end + 1;
if end <= self.text_buf.len() {
let json_slice = &self.text_buf[node.start..end];
let mapped: Vec<ParsedOrUnknown<T>> = deserialize_stream_map::<T>(json_slice);
if mapped.is_empty() {
items.push(StreamItem::Text(TextContent { text: json_slice.to_string() }));
} else {
let mut any_parsed = false;
for item in mapped {
match item {
ParsedOrUnknown::Parsed(v) => {
any_parsed = true;
items.push(StreamItem::Data(v));
}
ParsedOrUnknown::Unknown(u) => {
let u_end = u.end + 1;
if u_end <= json_slice.len() && u.start < u_end {
let sub = &json_slice[u.start..u_end];
items.push(StreamItem::Text(TextContent { text: sub.to_string() }));
} else {
debug!(target = "semantic_query::json_stream", "Skipping invalid unknown coordinates");
}
}
}
}
if !any_parsed {
items.push(StreamItem::Text(TextContent { text: json_slice.to_string() }));
}
}
self.last_offset = end;
}
}
items
}
pub fn finalize(self) -> Vec<StreamItem<T>> {
let mut items = Vec::new();
if self.last_offset < self.text_buf.len() {
let text_slice = &self.text_buf[self.last_offset..];
if !text_slice.trim().is_empty() {
items.push(StreamItem::Text(TextContent { text: text_slice.to_string() }));
}
}
items
}
pub fn text_buffer(&self) -> &str {
&self.text_buf
}
}
impl<T> Default for JsonStreamProcessor<T>
where
T: DeserializeOwned + JsonSchema + Send + 'static,
{
fn default() -> Self {
Self::new()
}
}
pub fn process_complete_text<T>(text: &str) -> Vec<StreamItem<T>>
where
T: DeserializeOwned + JsonSchema,
{
let mut items = Vec::new();
let roots = find_json_structures(text);
let mut cursor = 0usize;
for node in roots {
if node.start > cursor {
let text_slice = &text[cursor..node.start];
let trimmed = text_slice.trim();
if !trimmed.is_empty() {
items.push(StreamItem::Text(TextContent { text: text_slice.to_string() }));
}
}
let end = node.end + 1; let json_slice = &text[node.start..end];
let mapped: Vec<ParsedOrUnknown<T>> = deserialize_stream_map::<T>(json_slice);
if mapped.is_empty() {
items.push(StreamItem::Text(TextContent { text: json_slice.to_string() }));
} else {
let mut any_parsed = false;
for item in mapped {
match item {
ParsedOrUnknown::Parsed(v) => {
any_parsed = true;
items.push(StreamItem::Data(v));
}
ParsedOrUnknown::Unknown(u) => {
let u_end = u.end + 1;
if u_end <= json_slice.len() && u.start < u_end {
let sub = &json_slice[u.start..u_end];
items.push(StreamItem::Text(TextContent { text: sub.to_string() }));
} else {
debug!(target = "semantic_query::json_stream", "Skipping invalid unknown coordinates");
}
}
}
}
if !any_parsed {
items.push(StreamItem::Text(TextContent { text: json_slice.to_string() }));
}
}
cursor = end;
}
if cursor < text.len() {
let text_slice = &text[cursor..];
let trimmed = text_slice.trim();
if !trimmed.is_empty() {
items.push(StreamItem::Text(TextContent { text: text_slice.to_string() }));
}
}
items
}