Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T04:32:27.744Z Public web read
NIP-34 coordinate30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omega
MaintainersHidden in public view
References2 branches · 1 tag
Read-only clonegit clone https://openagents.com/git/tenant.openagents/omega.git
Browse files

streaming_parser.rs

1305 lines · 42.6 KB · rust
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
Served at tenant.openagents/omega Member data and write actions are omitted.