Skip to repository content1305 lines · 42.6 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T05:36:01.427Z Public web read
NIP-34 coordinate
30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omegaMaintainersHidden in public view
References2 branches · 1 tag
Read-only clone
git clone https://openagents.com/git/tenant.openagents/omega.gitBrowse files
streaming_parser.rs
1use smallvec::SmallVec;
2
3use super::{Edit, PartialEdit};
4
5/// Events emitted by `StreamingParser` for edit-mode input.
6#[derive(Debug, PartialEq, Eq)]
7pub enum EditEvent {
8 /// A chunk of `old_text` for an edit operation.
9 OldTextChunk {
10 edit_index: usize,
11 chunk: String,
12 done: bool,
13 },
14 /// A chunk of `new_text` for an edit operation.
15 NewTextChunk {
16 edit_index: usize,
17 chunk: String,
18 done: bool,
19 },
20}
21
22/// Events emitted by `StreamingParser` for write-mode input.
23#[derive(Debug, PartialEq, Eq)]
24pub enum WriteEvent {
25 /// A chunk of content for write/overwrite mode.
26 ContentChunk { chunk: String },
27}
28
29/// Tracks the streaming state of a single edit to detect deltas.
30#[derive(Default, Debug)]
31struct EditStreamState {
32 old_text_emitted_len: usize,
33 old_text_done: bool,
34 new_text_emitted_len: usize,
35 new_text_done: bool,
36 hold_until_complete: bool,
37 buffer_new_text_until_old_text_done: bool,
38}
39
40/// Converts incrementally-growing tool call JSON into a stream of chunk events.
41///
42/// The tool call streaming infrastructure delivers partial JSON objects where
43/// string fields grow over time. This parser compares consecutive partials,
44/// computes the deltas, and emits `EditEvent`s or `WriteEvent`s that downstream
45/// pipeline stages (`StreamingFuzzyMatcher` for old_text, `StreamingDiff` for
46/// new_text) can consume incrementally.
47///
48/// Because partial JSON comes through a fixer (`partial-json-fixer`) that
49/// closes incomplete escape sequences, a string can temporarily contain wrong
50/// trailing characters (e.g. a literal `\` instead of `\n`). We handle this
51/// by holding back trailing backslash characters in non-finalized chunks: if
52/// a partial string ends with `\` (0x5C), that byte is not emitted until the
53/// next partial confirms or corrects it. This avoids feeding corrupted bytes
54/// to downstream consumers.
55#[derive(Default, Debug)]
56pub struct StreamingParser {
57 edit_states: Vec<EditStreamState>,
58 content_emitted_len: usize,
59}
60
61impl StreamingParser {
62 /// Push a new set of partial edits (from edit mode) and return any events.
63 ///
64 /// Each call should pass the *entire current* edits array as seen in the
65 /// latest partial input. The parser will diff it against its internal state
66 /// to produce only the new events.
67 pub fn push_edits(&mut self, edits: &[PartialEdit]) -> SmallVec<[EditEvent; 4]> {
68 let mut events = SmallVec::new();
69
70 for (index, partial) in edits.iter().enumerate() {
71 if index >= self.edit_states.len() {
72 // A new edit appeared — finalize the previous one if there was one.
73 if let Some(previous) = self.finalize_previous_edit(
74 index,
75 edits
76 .get(index.saturating_sub(1))
77 .and_then(|edit| edit.old_text.as_deref()),
78 edits
79 .get(index.saturating_sub(1))
80 .and_then(|edit| edit.new_text.as_deref()),
81 ) {
82 events.extend(previous);
83 }
84 self.edit_states.push(EditStreamState::default());
85 }
86
87 let state = &mut self.edit_states[index];
88
89 if state.old_text_emitted_len == 0
90 && state.new_text_emitted_len == 0
91 && !state.old_text_done
92 && partial.new_text.is_some()
93 && !state.buffer_new_text_until_old_text_done
94 {
95 if partial
96 .old_text
97 .as_ref()
98 .is_some_and(|old_text| !old_text.is_empty())
99 {
100 state.hold_until_complete = true;
101 } else {
102 state.buffer_new_text_until_old_text_done = true;
103 }
104 }
105
106 if state.hold_until_complete {
107 continue;
108 }
109
110 // Process old_text changes.
111 if let Some(old_text) = &partial.old_text
112 && !state.old_text_done
113 {
114 if partial.new_text.is_some() && !state.buffer_new_text_until_old_text_done {
115 // new_text appeared after old_text, so old_text is done — emit everything.
116 let start = find_char_boundary(old_text, state.old_text_emitted_len);
117 let chunk = old_text[start..].to_string();
118 state.old_text_done = true;
119 state.old_text_emitted_len = old_text.len();
120 events.push(EditEvent::OldTextChunk {
121 edit_index: index,
122 chunk,
123 done: true,
124 });
125 } else {
126 let safe_end = safe_emit_end_for_edit_text(old_text);
127 let safe_start = find_char_boundary(old_text, state.old_text_emitted_len);
128
129 if safe_end > safe_start {
130 let chunk = old_text[safe_start..safe_end].to_string();
131 state.old_text_emitted_len = safe_end;
132 events.push(EditEvent::OldTextChunk {
133 edit_index: index,
134 chunk,
135 done: false,
136 });
137 }
138 }
139 }
140
141 // Process new_text changes.
142 if let Some(new_text) = &partial.new_text
143 && state.old_text_done
144 && !state.new_text_done
145 {
146 let safe_end = safe_emit_end_for_edit_text(new_text);
147 let safe_start = find_char_boundary(new_text, state.new_text_emitted_len);
148
149 if safe_end > safe_start {
150 let chunk = new_text[safe_start..safe_end].to_string();
151 state.new_text_emitted_len = safe_end;
152 events.push(EditEvent::NewTextChunk {
153 edit_index: index,
154 chunk,
155 done: false,
156 });
157 }
158 }
159 }
160
161 events
162 }
163
164 /// Push new content and return any events.
165 ///
166 /// Each call should pass the *entire current* content string. The parser
167 /// will diff it against its internal state to emit only the new chunk.
168 pub fn push_content(&mut self, content: &str) -> SmallVec<[WriteEvent; 1]> {
169 let mut events = SmallVec::new();
170
171 let safe_end = safe_emit_end(content);
172 let safe_start = find_char_boundary(content, self.content_emitted_len);
173 if safe_end > safe_start {
174 let chunk = content[safe_start..safe_end].to_string();
175 self.content_emitted_len = safe_end;
176 events.push(WriteEvent::ContentChunk { chunk });
177 }
178
179 events
180 }
181
182 /// Finalize all edits with the complete input. This emits `done: true`
183 /// events for any in-progress old_text or new_text that hasn't been
184 /// finalized yet.
185 ///
186 /// `final_edits` should be the fully deserialized final edits array. The
187 /// parser compares against its tracked state and emits any remaining deltas
188 /// with `done: true`.
189 pub fn finalize_edits(&mut self, edits: &[Edit]) -> SmallVec<[EditEvent; 4]> {
190 let mut events = SmallVec::new();
191
192 for (index, edit) in edits.iter().enumerate() {
193 if index >= self.edit_states.len() {
194 // This edit was never seen in partials — emit it fully.
195 if let Some(previous) = self.finalize_previous_edit(
196 index,
197 edits
198 .get(index.saturating_sub(1))
199 .map(|edit| edit.old_text.as_str()),
200 edits
201 .get(index.saturating_sub(1))
202 .map(|edit| edit.new_text.as_str()),
203 ) {
204 events.extend(previous);
205 }
206 self.edit_states.push(EditStreamState::default());
207 }
208
209 let state = &mut self.edit_states[index];
210
211 if state.hold_until_complete {
212 state.old_text_done = true;
213 state.old_text_emitted_len = edit.old_text.len();
214 state.new_text_done = true;
215 state.new_text_emitted_len = edit.new_text.len();
216 state.hold_until_complete = false;
217 state.buffer_new_text_until_old_text_done = false;
218 events.push(EditEvent::OldTextChunk {
219 edit_index: index,
220 chunk: edit.old_text.clone(),
221 done: true,
222 });
223 events.push(EditEvent::NewTextChunk {
224 edit_index: index,
225 chunk: edit.new_text.clone(),
226 done: true,
227 });
228 continue;
229 }
230
231 if !state.old_text_done {
232 let start = find_char_boundary(&edit.old_text, state.old_text_emitted_len);
233 let chunk = edit.old_text[start..].to_string();
234 state.old_text_done = true;
235 state.old_text_emitted_len = edit.old_text.len();
236 events.push(EditEvent::OldTextChunk {
237 edit_index: index,
238 chunk,
239 done: true,
240 });
241 }
242
243 if !state.new_text_done {
244 let start = find_char_boundary(&edit.new_text, state.new_text_emitted_len);
245 let chunk = edit.new_text[start..].to_string();
246 state.new_text_done = true;
247 state.new_text_emitted_len = edit.new_text.len();
248 events.push(EditEvent::NewTextChunk {
249 edit_index: index,
250 chunk,
251 done: true,
252 });
253 }
254 }
255
256 events
257 }
258
259 /// Finalize content with the complete input.
260 pub fn finalize_content(&mut self, content: &str) -> SmallVec<[WriteEvent; 1]> {
261 let mut events = SmallVec::new();
262
263 let start = find_char_boundary(content, self.content_emitted_len);
264 if content.len() > start {
265 let chunk = content[start..].to_string();
266 self.content_emitted_len = content.len();
267 events.push(WriteEvent::ContentChunk { chunk });
268 }
269
270 events
271 }
272
273 /// When a new edit appears at `index`, finalize the edit at `index - 1`
274 /// by emitting a `NewTextChunk { done: true }` if it hasn't been finalized.
275 fn finalize_previous_edit(
276 &mut self,
277 new_index: usize,
278 old_text: Option<&str>,
279 new_text: Option<&str>,
280 ) -> Option<SmallVec<[EditEvent; 2]>> {
281 if new_index == 0 || self.edit_states.is_empty() {
282 return None;
283 }
284
285 let previous_index = new_index - 1;
286 if previous_index >= self.edit_states.len() {
287 return None;
288 }
289
290 let state = &mut self.edit_states[previous_index];
291 let mut events = SmallVec::new();
292
293 if state.hold_until_complete {
294 let old_text = old_text.unwrap_or_default();
295 let new_text = new_text.unwrap_or_default();
296 state.old_text_done = true;
297 state.old_text_emitted_len = old_text.len();
298 state.new_text_done = true;
299 state.new_text_emitted_len = new_text.len();
300 state.hold_until_complete = false;
301 state.buffer_new_text_until_old_text_done = false;
302 events.push(EditEvent::OldTextChunk {
303 edit_index: previous_index,
304 chunk: old_text.to_string(),
305 done: true,
306 });
307 events.push(EditEvent::NewTextChunk {
308 edit_index: previous_index,
309 chunk: new_text.to_string(),
310 done: true,
311 });
312 return Some(events);
313 }
314
315 if !state.old_text_done {
316 let old_text = old_text.unwrap_or_default();
317 let start = find_char_boundary(old_text, state.old_text_emitted_len);
318 state.old_text_done = true;
319 state.old_text_emitted_len = old_text.len();
320 events.push(EditEvent::OldTextChunk {
321 edit_index: previous_index,
322 chunk: old_text[start..].to_string(),
323 done: true,
324 });
325 }
326
327 if !state.new_text_done {
328 let new_text = new_text.unwrap_or_default();
329 let start = find_char_boundary(new_text, state.new_text_emitted_len);
330 state.new_text_done = true;
331 state.new_text_emitted_len = new_text.len();
332 state.buffer_new_text_until_old_text_done = false;
333 events.push(EditEvent::NewTextChunk {
334 edit_index: previous_index,
335 chunk: new_text[start..].to_string(),
336 done: true,
337 });
338 }
339
340 Some(events)
341 }
342}
343
344/// Returns the byte position up to which it is safe to emit from a partial
345/// string. If the string ends with a backslash (`\`, 0x5C), that byte is
346/// held back because it may be an artifact of the partial JSON fixer closing
347/// an incomplete escape sequence (e.g. turning a half-received `\n` into `\\`).
348/// The next partial will reveal the correct character.
349///
350/// The returned position is always a valid UTF-8 character boundary.
351fn safe_emit_end(text: &str) -> usize {
352 if text.ends_with('\\') {
353 text.len() - 1
354 } else {
355 text.len()
356 }
357}
358
359fn safe_emit_end_for_edit_text(text: &str) -> usize {
360 let safe_end = safe_emit_end(text);
361 // Use string slicing to check the last character, ensuring we respect UTF-8 boundaries.
362 if safe_end > 0 && text[..safe_end].ends_with('\n') {
363 safe_end - 1
364 } else {
365 safe_end
366 }
367}
368
369/// Finds a valid UTF-8 character boundary at or before the target position.
370///
371/// When streaming partial JSON, the text structure can change between updates
372/// (e.g., an escape sequence being completed). This means a byte position that
373/// was valid in one partial may land inside a multi-byte character in the next.
374/// This function finds the nearest valid boundary at or before the target.
375fn find_char_boundary(text: &str, target: usize) -> usize {
376 if target >= text.len() {
377 return text.len();
378 }
379 if text.is_char_boundary(target) {
380 return target;
381 }
382 // Walk backwards to find a valid boundary.
383 let mut pos = target;
384 while pos > 0 && !text.is_char_boundary(pos) {
385 pos -= 1;
386 }
387 pos
388}
389
390#[cfg(test)]
391mod tests {
392 use super::*;
393 use proptest::prelude::*;
394
395 fn emitted_len_inside_multibyte_char() -> impl Strategy<Value = (String, String)> {
396 (1usize..8, prop::sample::select(&["。", "—", "é", "🦀"])).prop_map(
397 |(emitted_len, multibyte_char)| {
398 let first = "a".repeat(emitted_len);
399 let second = format!("{}{}", "a".repeat(emitted_len - 1), multibyte_char);
400 (first, second)
401 },
402 )
403 }
404
405 fn boundary_sensitive_text() -> impl Strategy<Value = String> {
406 prop_oneof![
407 emitted_len_inside_multibyte_char().prop_map(|(first, _)| first),
408 emitted_len_inside_multibyte_char().prop_map(|(_, second)| second),
409 prop::sample::select(&[
410 "",
411 "a",
412 "ab",
413 "ab\\",
414 "a。",
415 "a—",
416 "hello,\\",
417 "hello,\n",
418 "hello,\nworld",
419 ])
420 .prop_map(ToString::to_string),
421 ]
422 }
423
424 fn partial_edit() -> impl Strategy<Value = PartialEdit> {
425 (
426 prop::option::of(boundary_sensitive_text()),
427 prop::option::of(boundary_sensitive_text()),
428 )
429 .prop_map(|(old_text, new_text)| PartialEdit { old_text, new_text })
430 }
431
432 #[test]
433 fn test_first_edit_with_new_text_in_first_chunk_is_held_until_finalize() {
434 let mut parser = StreamingParser::default();
435
436 let events = parser.push_edits(&[PartialEdit {
437 old_text: Some("old".into()),
438 new_text: Some("new".into()),
439 }]);
440 assert!(events.is_empty());
441
442 let events = parser.push_edits(&[PartialEdit {
443 old_text: Some("old text".into()),
444 new_text: Some("new text".into()),
445 }]);
446 assert!(events.is_empty());
447
448 let events = parser.finalize_edits(&[Edit {
449 old_text: "old text".into(),
450 new_text: "new text".into(),
451 }]);
452 assert_eq!(
453 events.as_slice(),
454 &[
455 EditEvent::OldTextChunk {
456 edit_index: 0,
457 chunk: "old text".into(),
458 done: true,
459 },
460 EditEvent::NewTextChunk {
461 edit_index: 0,
462 chunk: "new text".into(),
463 done: true,
464 },
465 ]
466 );
467 }
468
469 #[test]
470 fn test_single_edit_streamed_incrementally() {
471 let mut parser = StreamingParser::default();
472
473 // old_text arrives in chunks: "hell" → "hello w" → "hello world"
474 let events = parser.push_edits(&[PartialEdit {
475 old_text: Some("hell".into()),
476 new_text: None,
477 }]);
478 assert_eq!(
479 events.as_slice(),
480 &[EditEvent::OldTextChunk {
481 edit_index: 0,
482 chunk: "hell".into(),
483 done: false,
484 }]
485 );
486
487 let events = parser.push_edits(&[PartialEdit {
488 old_text: Some("hello w".into()),
489 new_text: None,
490 }]);
491 assert_eq!(
492 events.as_slice(),
493 &[EditEvent::OldTextChunk {
494 edit_index: 0,
495 chunk: "o w".into(),
496 done: false,
497 }]
498 );
499
500 // new_text appears → old_text finalizes
501 let events = parser.push_edits(&[PartialEdit {
502 old_text: Some("hello world".into()),
503 new_text: Some("good".into()),
504 }]);
505 assert_eq!(
506 events.as_slice(),
507 &[
508 EditEvent::OldTextChunk {
509 edit_index: 0,
510 chunk: "orld".into(),
511 done: true,
512 },
513 EditEvent::NewTextChunk {
514 edit_index: 0,
515 chunk: "good".into(),
516 done: false,
517 },
518 ]
519 );
520
521 // new_text grows
522 let events = parser.push_edits(&[PartialEdit {
523 old_text: Some("hello world".into()),
524 new_text: Some("goodbye world".into()),
525 }]);
526 assert_eq!(
527 events.as_slice(),
528 &[EditEvent::NewTextChunk {
529 edit_index: 0,
530 chunk: "bye world".into(),
531 done: false,
532 }]
533 );
534
535 // Finalize
536 let events = parser.finalize_edits(&[Edit {
537 old_text: "hello world".into(),
538 new_text: "goodbye world".into(),
539 }]);
540 assert_eq!(
541 events.as_slice(),
542 &[EditEvent::NewTextChunk {
543 edit_index: 0,
544 chunk: "".into(),
545 done: true,
546 }]
547 );
548 }
549
550 #[test]
551 fn test_done_chunks_preserve_trailing_newlines() {
552 let mut parser = StreamingParser::default();
553
554 let events = parser.finalize_edits(&[Edit {
555 old_text: "before\n".into(),
556 new_text: "after\n".into(),
557 }]);
558 assert_eq!(
559 events.as_slice(),
560 &[
561 EditEvent::OldTextChunk {
562 edit_index: 0,
563 chunk: "before\n".into(),
564 done: true,
565 },
566 EditEvent::NewTextChunk {
567 edit_index: 0,
568 chunk: "after\n".into(),
569 done: true,
570 },
571 ]
572 );
573 }
574
575 #[test]
576 fn test_partial_edit_preserves_trailing_newlines() {
577 let mut parser = StreamingParser::default();
578
579 let events = parser.push_edits(&[PartialEdit {
580 old_text: Some("before\n".into()),
581 new_text: Some("after\n".into()),
582 }]);
583 assert!(events.is_empty());
584
585 let events = parser.finalize_edits(&[Edit {
586 old_text: "before\n".into(),
587 new_text: "after\n".into(),
588 }]);
589 assert_eq!(
590 events.as_slice(),
591 &[
592 EditEvent::OldTextChunk {
593 edit_index: 0,
594 chunk: "before\n".into(),
595 done: true,
596 },
597 EditEvent::NewTextChunk {
598 edit_index: 0,
599 chunk: "after\n".into(),
600 done: true,
601 },
602 ]
603 );
604 }
605
606 #[test]
607 fn test_multiple_edits_sequential() {
608 let mut parser = StreamingParser::default();
609
610 // First edit streams in
611 let events = parser.push_edits(&[PartialEdit {
612 old_text: Some("first old".into()),
613 new_text: None,
614 }]);
615 assert_eq!(
616 events.as_slice(),
617 &[EditEvent::OldTextChunk {
618 edit_index: 0,
619 chunk: "first old".into(),
620 done: false,
621 }]
622 );
623
624 let events = parser.push_edits(&[PartialEdit {
625 old_text: Some("first old".into()),
626 new_text: Some("first new".into()),
627 }]);
628 assert_eq!(
629 events.as_slice(),
630 &[
631 EditEvent::OldTextChunk {
632 edit_index: 0,
633 chunk: "".into(),
634 done: true,
635 },
636 EditEvent::NewTextChunk {
637 edit_index: 0,
638 chunk: "first new".into(),
639 done: false,
640 },
641 ]
642 );
643
644 // Second edit appears → first edit's new_text is finalized
645 let events = parser.push_edits(&[
646 PartialEdit {
647 old_text: Some("first old".into()),
648 new_text: Some("first new".into()),
649 },
650 PartialEdit {
651 old_text: Some("second".into()),
652 new_text: None,
653 },
654 ]);
655 assert_eq!(
656 events.as_slice(),
657 &[
658 EditEvent::NewTextChunk {
659 edit_index: 0,
660 chunk: "".into(),
661 done: true,
662 },
663 EditEvent::OldTextChunk {
664 edit_index: 1,
665 chunk: "second".into(),
666 done: false,
667 },
668 ]
669 );
670
671 // Finalize everything
672 let events = parser.finalize_edits(&[
673 Edit {
674 old_text: "first old".into(),
675 new_text: "first new".into(),
676 },
677 Edit {
678 old_text: "second old".into(),
679 new_text: "second new".into(),
680 },
681 ]);
682 assert_eq!(
683 events.as_slice(),
684 &[
685 EditEvent::OldTextChunk {
686 edit_index: 1,
687 chunk: " old".into(),
688 done: true,
689 },
690 EditEvent::NewTextChunk {
691 edit_index: 1,
692 chunk: "second new".into(),
693 done: true,
694 },
695 ]
696 );
697 }
698
699 #[test]
700 fn test_content_streamed_incrementally() {
701 let mut parser = StreamingParser::default();
702
703 let events = parser.push_content("hello");
704 assert_eq!(
705 events.as_slice(),
706 &[WriteEvent::ContentChunk {
707 chunk: "hello".into(),
708 }]
709 );
710
711 let events = parser.push_content("hello world");
712 assert_eq!(
713 events.as_slice(),
714 &[WriteEvent::ContentChunk {
715 chunk: " world".into(),
716 }]
717 );
718
719 // No change
720 let events = parser.push_content("hello world");
721 assert!(events.is_empty());
722
723 let events = parser.push_content("hello world!");
724 assert_eq!(
725 events.as_slice(),
726 &[WriteEvent::ContentChunk { chunk: "!".into() }]
727 );
728
729 // Finalize with no additional content
730 let events = parser.finalize_content("hello world!");
731 assert!(events.is_empty());
732 }
733
734 #[test]
735 fn test_finalize_content_with_remaining() {
736 let mut parser = StreamingParser::default();
737
738 parser.push_content("partial");
739 let events = parser.finalize_content("partial content here");
740 assert_eq!(
741 events.as_slice(),
742 &[WriteEvent::ContentChunk {
743 chunk: " content here".into(),
744 }]
745 );
746 }
747
748 #[test]
749 fn test_content_trailing_backslash_held_back() {
750 let mut parser = StreamingParser::default();
751
752 // Partial JSON fixer turns incomplete \n into \\ (literal backslash).
753 // The trailing backslash is held back.
754 let events = parser.push_content("hello,\\");
755 assert_eq!(
756 events.as_slice(),
757 &[WriteEvent::ContentChunk {
758 chunk: "hello,".into(),
759 }]
760 );
761
762 // Next partial corrects the escape to an actual newline.
763 // The held-back byte was wrong; the correct newline is emitted.
764 let events = parser.push_content("hello,\n");
765 assert_eq!(
766 events.as_slice(),
767 &[WriteEvent::ContentChunk { chunk: "\n".into() }]
768 );
769
770 // Normal growth.
771 let events = parser.push_content("hello,\nworld");
772 assert_eq!(
773 events.as_slice(),
774 &[WriteEvent::ContentChunk {
775 chunk: "world".into(),
776 }]
777 );
778 }
779
780 #[test]
781 fn test_content_finalize_with_trailing_backslash() {
782 let mut parser = StreamingParser::default();
783
784 // Stream a partial with a fixer-corrupted trailing backslash.
785 // The backslash is held back.
786 parser.push_content("abc\\");
787
788 // Finalize reveals the correct character.
789 let events = parser.finalize_content("abc\n");
790 assert_eq!(
791 events.as_slice(),
792 &[WriteEvent::ContentChunk { chunk: "\n".into() }]
793 );
794 }
795
796 proptest! {
797 #[test]
798 fn test_content_finalize_does_not_panic_when_emitted_len_lands_inside_multibyte_char(
799 pair in emitted_len_inside_multibyte_char()
800 ) {
801 let (first, second) = pair;
802 let mut parser = StreamingParser::default();
803
804 parser.push_content(&first);
805 parser.finalize_content(&second);
806 }
807
808 #[test]
809 fn test_push_edits_does_not_panic_on_boundary_sensitive_sequences(
810 partials in prop::collection::vec(prop::collection::vec(partial_edit(), 0..4), 1..12)
811 ) {
812 let mut parser = StreamingParser::default();
813
814 for edits in partials {
815 parser.push_edits(&edits);
816 }
817 }
818 }
819
820 #[test]
821 fn test_no_partials_direct_finalize() {
822 let mut parser = StreamingParser::default();
823
824 let events = parser.finalize_edits(&[Edit {
825 old_text: "old".into(),
826 new_text: "new".into(),
827 }]);
828 assert_eq!(
829 events.as_slice(),
830 &[
831 EditEvent::OldTextChunk {
832 edit_index: 0,
833 chunk: "old".into(),
834 done: true,
835 },
836 EditEvent::NewTextChunk {
837 edit_index: 0,
838 chunk: "new".into(),
839 done: true,
840 },
841 ]
842 );
843 }
844
845 #[test]
846 fn test_no_partials_direct_finalize_multiple() {
847 let mut parser = StreamingParser::default();
848
849 let events = parser.finalize_edits(&[
850 Edit {
851 old_text: "first old".into(),
852 new_text: "first new".into(),
853 },
854 Edit {
855 old_text: "second old".into(),
856 new_text: "second new".into(),
857 },
858 ]);
859 assert_eq!(
860 events.as_slice(),
861 &[
862 EditEvent::OldTextChunk {
863 edit_index: 0,
864 chunk: "first old".into(),
865 done: true,
866 },
867 EditEvent::NewTextChunk {
868 edit_index: 0,
869 chunk: "first new".into(),
870 done: true,
871 },
872 EditEvent::OldTextChunk {
873 edit_index: 1,
874 chunk: "second old".into(),
875 done: true,
876 },
877 EditEvent::NewTextChunk {
878 edit_index: 1,
879 chunk: "second new".into(),
880 done: true,
881 },
882 ]
883 );
884 }
885
886 #[test]
887 fn test_old_text_no_growth() {
888 let mut parser = StreamingParser::default();
889
890 let events = parser.push_edits(&[PartialEdit {
891 old_text: Some("same".into()),
892 new_text: None,
893 }]);
894 assert_eq!(
895 events.as_slice(),
896 &[EditEvent::OldTextChunk {
897 edit_index: 0,
898 chunk: "same".into(),
899 done: false,
900 }]
901 );
902
903 // Same old_text, no new_text → no events
904 let events = parser.push_edits(&[PartialEdit {
905 old_text: Some("same".into()),
906 new_text: None,
907 }]);
908 assert!(events.is_empty());
909 }
910
911 #[test]
912 fn test_old_text_none_then_appears() {
913 let mut parser = StreamingParser::default();
914
915 // Edit exists but old_text is None (field hasn't arrived yet)
916 let events = parser.push_edits(&[PartialEdit {
917 old_text: None,
918 new_text: None,
919 }]);
920 assert!(events.is_empty());
921
922 // old_text appears
923 let events = parser.push_edits(&[PartialEdit {
924 old_text: Some("text".into()),
925 new_text: None,
926 }]);
927 assert_eq!(
928 events.as_slice(),
929 &[EditEvent::OldTextChunk {
930 edit_index: 0,
931 chunk: "text".into(),
932 done: false,
933 }]
934 );
935 }
936
937 #[test]
938 fn test_new_text_before_old_text_buffers_new_text_but_streams_old_text() {
939 let mut parser = StreamingParser::default();
940
941 let events = parser.push_edits(&[PartialEdit {
942 old_text: None,
943 new_text: Some("new".into()),
944 }]);
945 assert!(events.is_empty());
946
947 let events = parser.push_edits(&[PartialEdit {
948 old_text: Some("old".into()),
949 new_text: Some("new".into()),
950 }]);
951 assert_eq!(
952 events.as_slice(),
953 &[EditEvent::OldTextChunk {
954 edit_index: 0,
955 chunk: "old".into(),
956 done: false,
957 }]
958 );
959
960 let events = parser.finalize_edits(&[Edit {
961 old_text: "old".into(),
962 new_text: "new".into(),
963 }]);
964 assert_eq!(
965 events.as_slice(),
966 &[
967 EditEvent::OldTextChunk {
968 edit_index: 0,
969 chunk: "".into(),
970 done: true,
971 },
972 EditEvent::NewTextChunk {
973 edit_index: 0,
974 chunk: "new".into(),
975 done: true,
976 },
977 ]
978 );
979 }
980
981 #[test]
982 fn test_three_edits_streamed() {
983 let mut parser = StreamingParser::default();
984
985 // Stream first edit
986 parser.push_edits(&[PartialEdit {
987 old_text: Some("a".into()),
988 new_text: Some("A".into()),
989 }]);
990
991 // Second edit appears
992 parser.push_edits(&[
993 PartialEdit {
994 old_text: Some("a".into()),
995 new_text: Some("A".into()),
996 },
997 PartialEdit {
998 old_text: Some("b".into()),
999 new_text: Some("B".into()),
1000 },
1001 ]);
1002
1003 // Third edit appears
1004 let events = parser.push_edits(&[
1005 PartialEdit {
1006 old_text: Some("a".into()),
1007 new_text: Some("A".into()),
1008 },
1009 PartialEdit {
1010 old_text: Some("b".into()),
1011 new_text: Some("B".into()),
1012 },
1013 PartialEdit {
1014 old_text: Some("c".into()),
1015 new_text: None,
1016 },
1017 ]);
1018
1019 assert_eq!(
1020 events.as_slice(),
1021 &[
1022 EditEvent::OldTextChunk {
1023 edit_index: 1,
1024 chunk: "b".into(),
1025 done: true,
1026 },
1027 EditEvent::NewTextChunk {
1028 edit_index: 1,
1029 chunk: "B".into(),
1030 done: true,
1031 },
1032 EditEvent::OldTextChunk {
1033 edit_index: 2,
1034 chunk: "c".into(),
1035 done: false,
1036 },
1037 ]
1038 );
1039
1040 // Finalize
1041 let events = parser.finalize_edits(&[
1042 Edit {
1043 old_text: "a".into(),
1044 new_text: "A".into(),
1045 },
1046 Edit {
1047 old_text: "b".into(),
1048 new_text: "B".into(),
1049 },
1050 Edit {
1051 old_text: "c".into(),
1052 new_text: "C".into(),
1053 },
1054 ]);
1055 assert_eq!(
1056 events.as_slice(),
1057 &[
1058 EditEvent::OldTextChunk {
1059 edit_index: 2,
1060 chunk: "".into(),
1061 done: true,
1062 },
1063 EditEvent::NewTextChunk {
1064 edit_index: 2,
1065 chunk: "C".into(),
1066 done: true,
1067 },
1068 ]
1069 );
1070 }
1071
1072 #[test]
1073 fn test_finalize_with_unseen_old_text() {
1074 let mut parser = StreamingParser::default();
1075
1076 // Only saw partial old_text, never saw new_text in partials
1077 parser.push_edits(&[PartialEdit {
1078 old_text: Some("partial".into()),
1079 new_text: None,
1080 }]);
1081
1082 let events = parser.finalize_edits(&[Edit {
1083 old_text: "partial old text".into(),
1084 new_text: "replacement".into(),
1085 }]);
1086 assert_eq!(
1087 events.as_slice(),
1088 &[
1089 EditEvent::OldTextChunk {
1090 edit_index: 0,
1091 chunk: " old text".into(),
1092 done: true,
1093 },
1094 EditEvent::NewTextChunk {
1095 edit_index: 0,
1096 chunk: "replacement".into(),
1097 done: true,
1098 },
1099 ]
1100 );
1101 }
1102
1103 #[test]
1104 fn test_repeated_pushes_with_no_change() {
1105 let mut parser = StreamingParser::default();
1106
1107 let events = parser.push_edits(&[PartialEdit {
1108 old_text: Some("stable".into()),
1109 new_text: None,
1110 }]);
1111 assert_eq!(
1112 events.as_slice(),
1113 &[EditEvent::OldTextChunk {
1114 edit_index: 0,
1115 chunk: "stable".into(),
1116 done: false,
1117 }]
1118 );
1119
1120 // Push the exact same data again
1121 let events = parser.push_edits(&[PartialEdit {
1122 old_text: Some("stable".into()),
1123 new_text: None,
1124 }]);
1125 assert!(events.is_empty());
1126
1127 // And again
1128 let events = parser.push_edits(&[PartialEdit {
1129 old_text: Some("stable".into()),
1130 new_text: None,
1131 }]);
1132 assert!(events.is_empty());
1133 }
1134
1135 #[test]
1136 fn test_old_text_trailing_backslash_held_back() {
1137 let mut parser = StreamingParser::default();
1138
1139 // Partial-json-fixer produces a literal backslash when the JSON stream
1140 // cuts in the middle of an escape sequence like \n. The parser holds
1141 // back the trailing backslash instead of emitting it.
1142 let events = parser.push_edits(&[PartialEdit {
1143 old_text: Some("hello,\\".into()), // fixer closed incomplete \n as \\
1144 new_text: None,
1145 }]);
1146 // The trailing `\` is held back — only "hello," is emitted.
1147 assert_eq!(
1148 events.as_slice(),
1149 &[EditEvent::OldTextChunk {
1150 edit_index: 0,
1151 chunk: "hello,".into(),
1152 done: false,
1153 }]
1154 );
1155
1156 // Next partial: the fixer corrects the escape to \n.
1157 // Because edit text also holds back a trailing newline, nothing new
1158 // is emitted yet.
1159 let events = parser.push_edits(&[PartialEdit {
1160 old_text: Some("hello,\n".into()),
1161 new_text: None,
1162 }]);
1163 assert!(events.is_empty());
1164
1165 // Continue normally. The held-back newline is emitted together with the
1166 // next content once it is no longer trailing.
1167 let events = parser.push_edits(&[PartialEdit {
1168 old_text: Some("hello,\nworld".into()),
1169 new_text: None,
1170 }]);
1171 assert_eq!(
1172 events.as_slice(),
1173 &[EditEvent::OldTextChunk {
1174 edit_index: 0,
1175 chunk: "\nworld".into(),
1176 done: false,
1177 }]
1178 );
1179 }
1180
1181 #[test]
1182 fn test_multiline_old_and_new_text() {
1183 let mut parser = StreamingParser::default();
1184
1185 let events = parser.push_edits(&[PartialEdit {
1186 old_text: Some("line1\nline2".into()),
1187 new_text: None,
1188 }]);
1189 assert_eq!(
1190 events.as_slice(),
1191 &[EditEvent::OldTextChunk {
1192 edit_index: 0,
1193 chunk: "line1\nline2".into(),
1194 done: false,
1195 }]
1196 );
1197
1198 let events = parser.push_edits(&[PartialEdit {
1199 old_text: Some("line1\nline2\nline3".into()),
1200 new_text: Some("LINE1\n".into()),
1201 }]);
1202 assert_eq!(
1203 events.as_slice(),
1204 &[
1205 EditEvent::OldTextChunk {
1206 edit_index: 0,
1207 chunk: "\nline3".into(),
1208 done: true,
1209 },
1210 EditEvent::NewTextChunk {
1211 edit_index: 0,
1212 chunk: "LINE1".into(),
1213 done: false,
1214 },
1215 ]
1216 );
1217
1218 let events = parser.push_edits(&[PartialEdit {
1219 old_text: Some("line1\nline2\nline3".into()),
1220 new_text: Some("LINE1\nLINE2\nLINE3".into()),
1221 }]);
1222 assert_eq!(
1223 events.as_slice(),
1224 &[EditEvent::NewTextChunk {
1225 edit_index: 0,
1226 chunk: "\nLINE2\nLINE3".into(),
1227 done: false,
1228 }]
1229 );
1230 }
1231
1232 #[test]
1233 fn test_multibyte_char_with_trailing_backslash() {
1234 // Reproduces a panic where the stored `old_text_emitted_len` from a previous
1235 // partial lands inside a multi-byte UTF-8 character in the current partial.
1236 //
1237 // Scenario: The JSON fixer produces a literal backslash when the stream cuts
1238 // mid-escape. If the *next* partial replaces that backslash with a multi-byte
1239 // character (e.g., em-dash '—'), the stored byte position is no longer valid.
1240 let mut parser = StreamingParser::default();
1241
1242 // First partial: text ends with backslash (held back by safe_emit_end).
1243 // "abc" = 3 bytes, backslash held back, so emitted_len = 3.
1244 let events = parser.push_edits(&[PartialEdit {
1245 old_text: Some("abc\\".into()),
1246 new_text: None,
1247 }]);
1248 assert_eq!(
1249 events.as_slice(),
1250 &[EditEvent::OldTextChunk {
1251 edit_index: 0,
1252 chunk: "abc".into(),
1253 done: false,
1254 }]
1255 );
1256
1257 // Second partial: the backslash is replaced by em-dash '—' (3 bytes: E2 80 94).
1258 // "ab—" = 2 + 3 = 5 bytes total, with em-dash at bytes 2..5.
1259 // The stored emitted_len (3) is inside the em-dash!
1260 // This should NOT panic.
1261 let events = parser.push_edits(&[PartialEdit {
1262 old_text: Some("ab—".into()),
1263 new_text: None,
1264 }]);
1265 // The parser should handle this gracefully.
1266 let _ = events;
1267 }
1268
1269 #[test]
1270 fn test_emitted_len_inside_multibyte_char_boundary() {
1271 // More direct reproduction: emitted_len points inside a multi-byte character.
1272 //
1273 // This can happen when:
1274 // 1. First partial has text where byte N is a valid boundary
1275 // 2. Second partial has *different* text where byte N is inside a multi-byte char
1276 let mut parser = StreamingParser::default();
1277
1278 // First partial: "ab" (2 bytes), backslash held back.
1279 // After processing: emitted_len = 2
1280 let events = parser.push_edits(&[PartialEdit {
1281 old_text: Some("ab\\".into()),
1282 new_text: None,
1283 }]);
1284 assert_eq!(
1285 events.as_slice(),
1286 &[EditEvent::OldTextChunk {
1287 edit_index: 0,
1288 chunk: "ab".into(),
1289 done: false,
1290 }]
1291 );
1292
1293 // Second partial: "a—" where em-dash starts at byte 1 and spans bytes 1-3.
1294 // Stored emitted_len = 2, but byte 2 is inside the em-dash!
1295 // This should NOT panic.
1296 let events = parser.push_edits(&[PartialEdit {
1297 old_text: Some("a—".into()),
1298 new_text: None,
1299 }]);
1300 // The parser should handle this gracefully.
1301 // We don't care exactly what it emits, just that it doesn't panic.
1302 let _ = events;
1303 }
1304}
1305