Skip to main content

trex/
streaming.rs

1//! Streaming scan: feed the input in chunks and recover exactly the
2//! matches a whole-input scan would produce.
3//!
4//! The contract is metamorphic: for any pattern and any way of cutting
5//! the input into chunks, the streaming result equals
6//! [`crate::scan`] on the concatenation. The scanner reaches that by
7//! committing only what cannot change once committed.
8//!
9//! ## What can and cannot be committed early
10//!
11//! A match commits early only when no later byte can affect it. Two
12//! things make a match depend on later bytes:
13//!
14//! - a quantifier or balanced group whose match could extend past the
15//!   chunk seam (`.*`, `\B(...)`), and
16//! - a content guard `~"lit"`, whose forward window can be satisfied by
17//!   a literal that has not arrived yet.
18//!
19//! So the scanner commits a prefix only up to a boundary where no token
20//! and no match straddles the cut, and it commits nothing early at all
21//! while the pattern carries a guard (the guard's forward dependence is
22//! unbounded). Whatever cannot be committed is retained and re-scanned
23//! when the next chunk arrives; `finish` scans the final remainder. The
24//! retained buffer is bounded by the longest unbroken run between safe
25//! boundaries, not by the whole input, except for guarded or
26//! boundary-spanning patterns, which retain until `finish` by
27//! necessity, not by shortcut.
28
29use crate::ast::Pattern;
30use crate::engine::Span;
31use crate::pattern_set::PatternSet;
32use crate::token::Token;
33
34/// What a stream scans: one pattern, or a set whose every match carries
35/// its member.
36enum Source {
37    One(Pattern),
38    Set(Box<PatternSet>),
39}
40
41/// A streaming scanner over one compiled pattern, or over a set.
42pub struct StreamScanner {
43    source: Source,
44    /// Unfinalized bytes: everything from `base` to the current end.
45    buf: Vec<u8>,
46    /// Absolute offset of `buf[0]` in the whole input.
47    base: usize,
48    /// Matches already committed, in absolute offsets, each with the set
49    /// member that made it, `0` for a stream over one pattern.
50    out: Vec<(usize, Span)>,
51    /// Whether the pattern depends on input beyond the lines its match spans,
52    /// which forbids any early commit and means the whole input must be
53    /// buffered. Every cut falls just after a newline, so a line anchor is
54    /// settled once its line has ended. See
55    /// [`Pattern::depends_on_more_than_its_lines`].
56    defers_commit: bool,
57    /// The most bytes retained at once so far.
58    peak: usize,
59    /// Whether the buffer's end is a boundary when it ends a line. The only
60    /// token that could span such a cut is a run of whitespace, so it is
61    /// unless the pattern reads whitespace itself.
62    tail: bool,
63    /// The most tokens a match can span, for a pattern whose length is
64    /// bounded: such a match commits once the decided prefix holds that
65    /// many tokens past its start, and only the last tokens of the prefix
66    /// can still begin one that reaches past it.
67    span_limit: Option<usize>,
68    /// The absolute end of the last committed match. A scan of the retained
69    /// buffer resumes there, so a committed match is never found twice
70    /// where the buffer keeps its bytes for the lexer's sake.
71    committed_through: usize,
72    /// The shapes the retained buffer is lexed under, which for one pattern
73    /// are the shapes its own scan would build: holding them here is what
74    /// lets a push lex once and hand the tokens to the scan rather than have
75    /// it lex the same bytes again. Empty for a set, whose members can need
76    /// different shapes from one another and take their own lex.
77    shapes: crate::custom::ShapeSet,
78    /// The shapes a caller declared for one pattern, without the library's,
79    /// which the scan at the stream's end is handed as a whole-input scan
80    /// under them would be.
81    declared: crate::custom::ShapeSet,
82    /// The lexer's buffers, held from one push to the next so their pages stay
83    /// mapped and the token vector is not rebuilt per chunk.
84    lex_ws: crate::parallel_lex::TokenWorkspace,
85    /// Every byte of the retained buffer a cut may fall on, ascending, from the
86    /// one scan a push makes of it. Held rather than rebuilt for the same
87    /// reason the token buffer is.
88    bounds: Vec<usize>,
89}
90
91impl StreamScanner {
92    /// Begin a streaming scan for `pattern`.
93    #[must_use]
94    pub fn new(pattern: Pattern) -> Self {
95        Self::with_shapes(pattern, crate::custom::ShapeSet::new())
96    }
97
98    /// Begin a streaming scan for `pattern` under the shapes and kinds
99    /// `shapes` declares, which decide token boundaries as they do for
100    /// [`crate::scan_with_shapes`]: the stream finds the matches that scan
101    /// finds over the whole input.
102    #[must_use]
103    pub fn with_shapes(pattern: Pattern, shapes: crate::custom::ShapeSet) -> Self {
104        let span_limit = pattern.max_tokens();
105        let defers_commit = pattern.depends_on_more_than_its_lines() || !settles(&pattern, span_limit);
106        let tail = !pattern.reads_whitespace();
107        let mut scanner = Self::over(Source::One(pattern), defers_commit, tail, span_limit);
108        if !shapes.is_empty() {
109            if let Source::One(pattern) = &scanner.source {
110                scanner.shapes = shapes.with_library_shapes(&pattern.library_kinds());
111            }
112            scanner.declared = shapes;
113        }
114        scanner
115    }
116
117    /// Begin a streaming scan for every member of `set` at once, each match
118    /// tagged with its member. The stream commits under its most
119    /// conservative member: a match commits once the decided prefix holds
120    /// the longest span any member can make, nothing commits before the
121    /// end where any member depends on the whole input, and a line's end
122    /// is no boundary where any member reads whitespace; one window serves
123    /// them all.
124    #[must_use]
125    pub fn over_set(set: PatternSet) -> Self {
126        let span_limit = set.max_tokens();
127        let defers_commit = set.depends_on_more_than_its_lines()
128            || !set.patterns().iter().all(|p| settles(p, span_limit));
129        let tail = !set.reads_whitespace();
130        Self::over(Source::Set(Box::new(set)), defers_commit, tail, span_limit)
131    }
132
133    fn over(source: Source, defers_commit: bool, tail: bool, span_limit: Option<usize>) -> Self {
134        let shapes = match &source {
135            Source::One(pattern) => {
136                crate::custom::ShapeSet::new().with_library_shapes(&pattern.library_kinds())
137            }
138            Source::Set(_) => crate::custom::ShapeSet::new(),
139        };
140        Self {
141            shapes,
142            declared: crate::custom::ShapeSet::new(),
143            lex_ws: crate::parallel_lex::TokenWorkspace::default(),
144            bounds: Vec::new(),
145            source,
146            buf: Vec::new(),
147            base: 0,
148            out: Vec::new(),
149            defers_commit,
150            peak: 0,
151            tail,
152            span_limit,
153            committed_through: 0,
154        }
155    }
156
157    /// Whether a match can commit before the stream ends. A pattern that
158    /// depends on the whole input retains everything until [`Self::finish`].
159    #[must_use]
160    pub fn commits_early(&self) -> bool {
161        !self.defers_commit
162    }
163
164    /// The absolute offset of the first retained byte: everything before it
165    /// is committed, and no later byte can reach it.
166    #[must_use]
167    pub fn base(&self) -> usize {
168        self.base
169    }
170
171    /// The bytes retained now.
172    #[must_use]
173    pub fn retained(&self) -> usize {
174        self.buf.len()
175    }
176
177    /// Count this stream's offsets from its first retained byte rather than
178    /// from where it began, and give how far they moved; a reader adds that
179    /// to every offset reported after. A span's offsets reach four GiB, so a
180    /// stream that runs longer, as a file followed for days can, is rebased
181    /// as it goes. Every committed match must have been drained first, since
182    /// those carry the offsets from before.
183    ///
184    /// # Panics
185    ///
186    /// A committed match has not been drained.
187    pub fn rebase(&mut self) -> usize {
188        assert!(self.out.is_empty(), "a stream is rebased only once its committed matches are drained");
189        let by = self.base;
190        self.base = 0;
191        self.committed_through = self.committed_through.saturating_sub(by);
192        by
193    }
194
195    /// The most bytes retained at once so far.
196    #[must_use]
197    pub fn peak_retained(&self) -> usize {
198        self.peak
199    }
200
201    /// Feed the next chunk of input. Matches that can no longer change
202    /// are committed; the rest is retained for the following chunk.
203    pub fn push(&mut self, chunk: &[u8]) {
204        let taking = crate::trace::phase("streaming: taking the chunk in");
205        self.buf.extend_from_slice(chunk);
206        drop(taking);
207        self.peak = self.peak.max(self.buf.len());
208        if self.defers_commit {
209            // A guard's forward window or a field's comma count reaches
210            // outside the committable prefix, so no match is final until
211            // `finish` sees the whole input.
212            return;
213        }
214        match self.span_limit {
215            Some(limit) => self.commit_bounded(limit),
216            None => self.try_commit(),
217        }
218    }
219
220    /// The matches of the retained buffer from the last committed match on,
221    /// each with its member.
222    fn scan_retained(&self) -> Vec<(usize, Span)> {
223        let from = self.committed_through.saturating_sub(self.base);
224        match &self.source {
225            Source::One(pattern) => crate::engine::scan_with_shapes_from(pattern, &self.buf, &self.declared, from)
226                .into_iter()
227                .map(|s| (0, s))
228                .collect(),
229            Source::Set(set) => set.scan_from(&self.buf, from),
230        }
231    }
232
233    /// Lex `buf` under `shapes` into `ws`, which for one pattern is the lex its
234    /// own scan would take, so the scan can be handed these rather than repeat
235    /// them.
236    ///
237    /// Serial, into a buffer held from the last push, which is the pairing
238    /// [`crate::lexer::lex_into`] documents: a held vector has its pages
239    /// already and lexes in three quarters of the time, but holding one across
240    /// dispatched work costs more than the fresh pages save, because writing it
241    /// again invalidates its lines in every core's cache. A push lexes and then
242    /// walks the buffer serially, so the held half applies and the dispatched
243    /// half does not.
244    ///
245    /// A declared shape decides token boundaries `lex_into` does not read, so
246    /// that case takes the shaped lexer and a fresh vector.
247    ///
248    /// Takes the three pieces rather than `self` so the buffer, the shapes and
249    /// the workspace are borrowed apart.
250    fn lex_retained_into(
251        buf: &[u8],
252        shapes: &crate::custom::ShapeSet,
253        ws: &mut crate::parallel_lex::TokenWorkspace,
254    ) {
255        if shapes.is_empty() {
256            crate::lexer::lex_into(buf, &mut ws.toks);
257        } else {
258            let blobs = crate::lexer::blob_runs(buf);
259            ws.toks = crate::lexer::lex_with_shapes(buf, &blobs, shapes, 0);
260        }
261    }
262
263    /// [`Self::scan_retained`] over a lex of the buffer the caller holds
264    /// already, where that lex is the one the scan would have taken itself.
265    ///
266    /// [`Self::lex_retained`] lexes under the shapes one pattern's scan builds
267    /// from its own library kinds, so the tokens agree for every such pattern.
268    /// A set is lent that same lex and spends it a member at a time: a member
269    /// drawing no shape of its own reads these tokens, and one that does lexes
270    /// under its shapes, since no single lex serves members that disagree.
271    fn scan_retained_over(&self, toks: &[Token]) -> Vec<(usize, Span)> {
272        let from = self.committed_through.saturating_sub(self.base);
273        match &self.source {
274            Source::One(pattern) => {
275                // The routes over the retained buffer's own bytes, which answer
276                // without walking its tokens where the pattern draws no shape.
277                // The lex above is paid whatever they say - the drain reads it
278                // to size what the buffer keeps - so what a route saves here is
279                // the walk and not the lex.
280                if self.shapes.is_empty()
281                    && let Some(spans) = crate::engine::routed_spans_at(pattern, &self.buf, from)
282                {
283                    return spans.into_iter().map(|s| (0, s)).collect();
284                }
285                crate::engine::scan_over_tokens_from(pattern, &self.buf, toks, from)
286                    .into_iter()
287                    .map(|s| (0, s))
288                    .collect()
289            }
290            Source::Set(set) => set.scan_from_over(&self.buf, from, toks),
291        }
292    }
293
294    /// Commit for a pattern whose match spans at most `limit` tokens. The
295    /// decided prefix runs to the last line boundary, or to the buffer's end
296    /// when it ends a line; a match commits when the prefix holds `limit`
297    /// tokens from its start, since every alternative at that start is then
298    /// in view; and after the last commit only the prefix's last `limit - 1`
299    /// tokens can still begin a match reaching past it, so the buffer keeps
300    /// from the line boundary before them.
301    fn commit_bounded(&mut self, limit: usize) {
302        // The four acts of a push are timed apart, because the three that are
303        // not the lex were inside no phase at all and a chunk size that moved
304        // the whole could not be attributed to any of them.
305        let deciding = crate::trace::phase("streaming: deciding the prefix");
306        let end_held = scan_boundaries(&self.buf, &mut self.bounds);
307        let decided = boundary_below(&self.buf, &self.bounds, end_held, self.buf.len(), self.tail);
308        drop(deciding);
309        let Some(decided) = decided else {
310            return;
311        };
312        // The same re-lex as the unbounded path does, in the other function: a
313        // pattern whose match spans a bounded number of tokens commits here and
314        // never reaches `settled_through`, so an instrument on one path alone
315        // reports nothing for half the patterns. The names are shared so the
316        // lex reads as one quantity whichever path paid it.
317        let lexing = crate::trace::phase("streaming: lexing the retained buffer");
318        Self::lex_retained_into(&self.buf, &self.shapes, &mut self.lex_ws);
319        let toks = &self.lex_ws.toks;
320        drop(lexing);
321        // What a push walks and re-walks: the retained bytes and the tokens made
322        // of them. Four of a push's acts are linear in one or the other, so
323        // these two summed over the pushes, against the pushes, are the working
324        // set a chunk size asks a core to hold - which is a count rather than a
325        // reading of any cache.
326        crate::trace::counted(
327            "streaming: bytes the buffer holds",
328            u64::try_from(self.buf.len()).expect("a buffer within the counter's width"),
329        );
330        crate::trace::counted(
331            "streaming: bytes the tokens occupy",
332            u64::try_from(toks.len() * std::mem::size_of::<Token>())
333                .expect("a token vector within the counter's width"),
334        );
335        // Where each significant token of the decided prefix begins, which is
336        // the whole of what the drain below asks of them: how many start before
337        // a byte, and where the k-th one starts. A token's end decides only
338        // whether it lies inside the prefix, so nothing past this filter reads
339        // one and the list carries starts alone.
340        let gathering = crate::trace::phase("streaming: gathering the significant tokens");
341        let sig: Vec<usize> = toks
342            .iter()
343            .filter(|t| t.is_significant() && t.end() <= decided)
344            .map(crate::token::Token::start)
345            .collect();
346        let index_at = |byte: usize| sig.partition_point(|&s| s < byte);
347        let mut after = index_at(self.committed_through.saturating_sub(self.base));
348        drop(gathering);
349        let scanning = crate::trace::phase("streaming: scanning the retained buffer");
350        // What the scan is over, so its cost can be divided by its bytes. A
351        // scan resuming at `committed_through` covers the buffer from there,
352        // and that offset does not advance while nothing commits - so this says
353        // whether the phase is many small scans or a few large ones without
354        // anyone having to reason about which.
355        crate::trace::counted(
356            "streaming: bytes the scans cover",
357            u64::try_from(self.buf.len() - self.committed_through.saturating_sub(self.base))
358                .expect("a buffer within the counter's width"),
359        );
360        let found = self.scan_retained_over(toks);
361        drop(scanning);
362        for (member, m) in found {
363            if m.end() > decided || index_at(m.start()) + limit > sig.len() {
364                break;
365            }
366            self.out.push((member, shifted(m, self.base)));
367            self.committed_through = self.base + m.end();
368            after = index_at(m.end());
369        }
370        let dropping = crate::trace::phase("streaming: dropping the settled bytes");
371        let retain = after.max((sig.len() + 1).saturating_sub(limit));
372        let retain_byte = if retain < sig.len() { sig[retain] } else { decided };
373        // The second limit answered from the first scan's boundaries rather
374        // than from a second scan of the same bytes.
375        let cut = if retain_byte >= self.buf.len() {
376            self.buf.len()
377        } else {
378            boundary_below(&self.buf, &self.bounds, end_held, retain_byte + 1, self.tail)
379                .unwrap_or(0)
380        };
381        // What the drain moves: the bytes it keeps, which it shifts down to the
382        // buffer's start. A cost that tracks this rather than the chunk is one
383        // that depends on how much of a buffer survives a push, which is the
384        // shape a curve non-monotonic in chunk size would have.
385        crate::trace::counted(
386            "streaming: bytes the drain moves",
387            u64::try_from(self.buf.len() - cut).expect("a buffer within the counter's width"),
388        );
389        self.buf.drain(..cut);
390        self.base += cut;
391        drop(dropping);
392    }
393
394    /// The matches committed so far, handed over once: a later call returns
395    /// only what was committed after it, and `finish` returns what it
396    /// commits plus whatever was never drained.
397    pub fn drain_committed(&mut self) -> Vec<Span> {
398        self.drain_committed_with_members().into_iter().map(|(_, s)| s).collect()
399    }
400
401    /// [`Self::drain_committed`], each match with the set member that made
402    /// it; `0` throughout for a stream over one pattern.
403    pub fn drain_committed_with_members(&mut self) -> Vec<(usize, Span)> {
404        std::mem::take(&mut self.out)
405    }
406
407    /// Finish the stream: scan the retained buffer and return every match
408    /// not yet drained, in absolute offsets.
409    #[must_use]
410    pub fn finish(self) -> Vec<Span> {
411        self.finish_with_members().into_iter().map(|(_, s)| s).collect()
412    }
413
414    /// [`Self::finish`], each match with the set member that made it.
415    #[must_use]
416    pub fn finish_with_members(mut self) -> Vec<(usize, Span)> {
417        for (member, m) in self.scan_retained() {
418            self.out.push((member, shifted(m, self.base)));
419        }
420        self.out
421    }
422
423    /// Commit every match that ends at or before the largest boundary no
424    /// live attempt can still reach back across, then drop the committed
425    /// bytes.
426    ///
427    /// This is the path for a pattern whose match has no bounded length, so
428    /// there is no count of tokens that puts a limit on how far back a match
429    /// arriving with the next chunk could begin. The matches found in the
430    /// buffer as it stands do not answer that: `\W+ "END"` over `alpha\n`
431    /// holds no match at all, and cutting on that basis drops `alpha` from a
432    /// match the next chunk completes. What answers it is the walk's own
433    /// thread list at the end of the buffer - every attempt a further token
434    /// could carry forward - and the earliest byte any of those began at is
435    /// the first byte this buffer must keep.
436    ///
437    /// Where that cannot be computed, which is where the pattern needs the
438    /// set-reachability engine, nothing is dropped: a stream that keeps too
439    /// much is slow, and one that drops a live attempt is wrong.
440    fn try_commit(&mut self) {
441        // One lex of the retained buffer serves both the settling walk and the
442        // scan that follows it, which read the same bytes under the same
443        // shapes.
444        let lexing = crate::trace::phase("streaming: lexing the retained buffer");
445        Self::lex_retained_into(&self.buf, &self.shapes, &mut self.lex_ws);
446        let toks = &self.lex_ws.toks;
447        drop(lexing);
448        let Some(settled) = self.settled_through(toks) else {
449            return;
450        };
451        let scanning = crate::trace::phase("streaming: scanning the retained buffer");
452        crate::trace::counted(
453            "streaming: bytes the scans cover",
454            u64::try_from(self.buf.len() - self.committed_through.saturating_sub(self.base))
455                .expect("a buffer within the counter's width"),
456        );
457        let matches = self.scan_retained_over(toks);
458        drop(scanning);
459        // One scan of the buffer's boundaries, which both the empty case below
460        // and the straddle loop under it read: the loop asks for the next
461        // boundary down once per match that crosses a cut, and each of those
462        // asks was a scan of the whole buffer.
463        // Held to the end of the call rather than dropped at a point, because
464        // every way out of this function from here drops the bytes it settled
465        // or returns having settled none.
466        let _dropping = crate::trace::phase("streaming: dropping the settled bytes");
467        let end_held = scan_boundaries(&self.buf, &mut self.bounds);
468        if matches.is_empty() {
469            // No match yet, and no attempt reaching back past `settled`, so
470            // the bytes before the last boundary under it are dead weight.
471            if let Some(cut) = boundary_below(&self.buf, &self.bounds, end_held, settled, false) {
472                self.buf.drain(..cut);
473                self.base += cut;
474            }
475            return;
476        }
477        // The largest such boundary that no match straddles either. A match
478        // straddling the cut would be split, so such a cut is rejected.
479        let spans: Vec<Span> = matches.iter().map(|&(_, s)| s).collect();
480        let Some(cut) = largest_uncrossed_boundary(
481            &self.buf,
482            &self.bounds,
483            end_held,
484            &spans,
485            settled,
486            false,
487        ) else {
488            return;
489        };
490        for &(member, m) in &matches {
491            if m.end() <= cut {
492                self.out.push((member, shifted(m, self.base)));
493            }
494        }
495        self.buf.drain(..cut);
496        self.base += cut;
497    }
498
499    /// The byte of the retained buffer before which no attempt is still
500    /// running, so nothing below it can be reached by a match the next chunk
501    /// completes; `None` where the engine cannot say and the buffer must be
502    /// kept whole.
503    ///
504    /// A set is settled only as far as its least settled member: one member
505    /// with a live attempt holds the buffer for all of them, which is the
506    /// same rule as the one window the set commits under.
507    fn settled_through(&self, toks: &[Token]) -> Option<usize> {
508        // The walk consults the absent guard literals itself, so a
509        // short-circuit above this one could save the lex and not the walk -
510        // which is why the caller holds the lex and this phase times the walk
511        // alone.
512        let _walking = crate::trace::phase("streaming: the walk over the tokens");
513        let of = |p: &Pattern| {
514            crate::nfa::earliest_unsettled(p, &self.buf, toks).map(|at| at.unwrap_or(self.buf.len()))
515        };
516        match &self.source {
517            Source::One(pattern) => of(pattern),
518            Source::Set(set) => {
519                let mut least = self.buf.len();
520                for p in set.patterns() {
521                    least = least.min(of(p)?);
522                }
523                Some(least)
524            }
525        }
526    }
527}
528
529/// Where a held stream's offsets are counted from once its retained window
530/// has moved this far: a span's offsets reach four GiB, so a stream longer
531/// than that, as a file followed for days can be, counts from its retained
532/// window rather than from its start and adds the difference back.
533const REBASE_AT: usize = 1 << 30;
534
535/// A stream scanner and the bytes from its retained window on, with where
536/// they stand in the input: what a reader needs beside the scanner to read
537/// each committed match's text and registers and to place it on its line.
538///
539/// The bytes a push settles are dropped at the next push, so the matches a
540/// push commits are read against [`HeldStream::held`] until then.
541pub struct HeldStream {
542    scanner: StreamScanner,
543    /// The input from the scanner's retained window on.
544    held: Vec<u8>,
545    /// Where `held` begins, in the scanner's offsets.
546    held_base: usize,
547    /// How far the input has been pushed to the scanner, in its offsets.
548    pushed: usize,
549    /// Where the scanner's offset zero stands in the input.
550    origin: usize,
551    /// The input's newlines before `held`, where the stream counts them.
552    lines_before: Option<usize>,
553}
554
555/// What a held stream comes to at its input's end: the matches the scanner
556/// held until then, as spans over `held`, with each member, and the bytes
557/// and their place, as [`HeldStream`] gives them.
558pub struct Ended {
559    pub matches: Vec<(usize, Span)>,
560    pub held: Vec<u8>,
561    /// Where `held` begins in the input.
562    pub base: usize,
563    /// How many of the input's lines come before `held`, where the stream
564    /// counted them.
565    pub lines_before: Option<usize>,
566}
567
568impl HeldStream {
569    /// `scanner` over an input whose first byte stands at `origin`, after
570    /// `lines_before` of its lines; the lines are counted on as the stream
571    /// moves where they were counted to its start, and left uncounted where
572    /// they were not.
573    #[must_use]
574    pub fn new(scanner: StreamScanner, origin: usize, lines_before: Option<usize>) -> Self {
575        HeldStream { scanner, held: Vec::new(), held_base: 0, pushed: 0, origin, lines_before }
576    }
577
578    /// Whether a match can commit before the input ends.
579    #[must_use]
580    pub fn commits_early(&self) -> bool {
581        self.scanner.commits_early()
582    }
583
584    /// The bytes held: the input from the scanner's retained window on.
585    #[must_use]
586    pub fn held(&self) -> &[u8] {
587        &self.held
588    }
589
590    /// Where the held bytes begin in the input.
591    #[must_use]
592    pub fn base(&self) -> usize {
593        self.origin + self.held_base
594    }
595
596    /// How many of the input's lines come before the held bytes, where the
597    /// stream counts them.
598    #[must_use]
599    pub fn lines_before(&self) -> Option<usize> {
600        self.lines_before
601    }
602
603    /// How much of the held bytes is final: no match can still begin before
604    /// this offset into them, and every match that ends before it has been
605    /// committed.
606    #[must_use]
607    pub fn settled(&self) -> usize {
608        self.scanner.base() - self.held_base
609    }
610
611    /// Push the input's next bytes: the matches they commit, each with the
612    /// member that made it, as spans over the held bytes.
613    pub fn push(&mut self, bytes: &[u8]) -> Vec<(usize, Span)> {
614        self.settle();
615        self.held.extend_from_slice(bytes);
616        let from = self.pushed - self.held_base;
617        self.scanner.push(&self.held[from..]);
618        self.pushed = self.held_base + self.held.len();
619        let committed = self.scanner.drain_committed_with_members();
620        local(committed, self.held_base)
621    }
622
623    /// End the input: the matches the scanner held until then, over the
624    /// bytes still held.
625    #[must_use]
626    pub fn finish(mut self) -> Ended {
627        self.settle();
628        let HeldStream { scanner, held, held_base, origin, lines_before, .. } = self;
629        let matches = local(scanner.finish_with_members(), held_base);
630        Ended { matches, held, base: origin + held_base, lines_before }
631    }
632
633    /// Drop the bytes the last push settled, counting their newlines, and
634    /// count the offsets from the retained window once it is far along.
635    fn settle(&mut self) {
636        let keep_from = self.scanner.base();
637        if keep_from > self.held_base {
638            let dropped = keep_from - self.held_base;
639            if let Some(lines) = self.lines_before.as_mut() {
640                *lines += crate::byte_simd::count_byte(&self.held[..dropped], b'\n');
641            }
642            self.held.drain(..dropped);
643            self.held_base = keep_from;
644        }
645        if self.held_base >= REBASE_AT {
646            let by = self.scanner.rebase();
647            self.origin += by;
648            self.held_base -= by;
649            self.pushed -= by;
650        }
651    }
652}
653
654/// `committed`, spans in a scanner's offsets, as spans over the bytes held
655/// from `held_base`.
656fn local(committed: Vec<(usize, Span)>, held_base: usize) -> Vec<(usize, Span)> {
657    let at = |offset: usize| {
658        u32::try_from(offset - held_base).expect("a retained window is narrower than a span's width")
659    };
660    committed.into_iter().map(|(member, s)| (member, Span { start: at(s.start()), end: at(s.end()) })).collect()
661}
662
663/// Whether a stream over `pattern` can ever drop a byte before its end.
664///
665/// A bounded match settles by its own token count, whatever engine walks
666/// it. An unbounded one settles only where the walk's thread list can be
667/// read, which is where the pattern compiles to the single-pass program; a
668/// pattern the set-reachability engine owns - a balanced group, a field
669/// node - offers no such list, so nothing about it is ever settled and the
670/// stream must keep every byte until it ends. Deciding that here rather
671/// than per push is what lets the scanner report it: a caller asking
672/// [`Self::commits_early`] is told before the bytes pile up, and the
673/// command line says how much it retained when the stream closes.
674fn settles(pattern: &Pattern, span_limit: Option<usize>) -> bool {
675    span_limit.is_some() || crate::nfa::earliest_unsettled(pattern, b"", &[]).is_some()
676}
677
678/// The span at the whole input's offsets: the buffer's offsets plus `base`.
679fn shifted(span: Span, base: usize) -> Span {
680    let at = |offset: usize| {
681        u32::try_from(offset + base)
682            .expect("a stream's absolute offset fits the span width as an input's does")
683    };
684    Span { start: at(span.start()), end: at(span.end()) }
685}
686
687/// Scan `chunks` as a stream and return the whole-input matches. A
688/// convenience over [`StreamScanner`] for callers that already hold an
689/// iterator of byte slices.
690#[must_use]
691pub fn scan_chunked<'a>(pattern: &Pattern, chunks: impl IntoIterator<Item = &'a [u8]>) -> Vec<Span> {
692    let mut s = StreamScanner::new(pattern.clone());
693    for c in chunks {
694        s.push(c);
695    }
696    s.finish()
697}
698
699/// The largest safe boundary in `buf` below `limit` that no match
700/// straddles, or `None` when none qualifies (so nothing can be committed
701/// yet).
702/// Each rejected boundary sends this back for the next one down, and reading
703/// that from the boundaries already collected makes the retry a search rather
704/// than another scan of the buffer.
705fn largest_uncrossed_boundary(
706    buf: &[u8],
707    bounds: &[usize],
708    end_held: bool,
709    matches: &[Span],
710    limit: usize,
711    tail: bool,
712) -> Option<usize> {
713    let mut cut = boundary_below(buf, bounds, end_held, limit, tail)?;
714    loop {
715        // A match straddles `cut` when it starts before and ends after.
716        if matches.iter().any(|m| m.start() < cut && m.end() > cut) {
717            cut = boundary_below(buf, bounds, end_held, cut, tail)?;
718            continue;
719        }
720        return Some(cut);
721    }
722}
723
724/// Whether a token of the lexer's may span the newline at `nl`, so that no cut
725/// falls after it: a string, of either quote, carried over the line by a
726/// backslash before the newline or before its CR, or a char literal holding
727/// the newline itself, a quote on each side. Every other newline is a token
728/// boundary - a string closes on its own line - and a cut after it needs no
729/// quote state carried from the buffer's start.
730///
731/// A quote before the newline and none after it yet, at the buffer's end, may
732/// be the first half of such a char literal, so it holds the newline too.
733fn newline_held(buf: &[u8], nl: usize) -> bool {
734    if nl == 0 {
735        return false;
736    }
737    match buf[nl - 1] {
738        b'\'' => nl + 1 >= buf.len() || crate::lexer::char_literal_end(buf, nl - 1) == Some(nl + 2),
739        _ => crate::lexer::backslash_before_newline(buf, nl),
740    }
741}
742
743/// Every byte of `buf` a cut may fall on, ascending, written to `out`, and
744/// whether a token may run on past the buffer's end across its last newline.
745///
746/// A scan per limit would find the same newlines once per limit. This finds
747/// them once and [`boundary_below`] answers each limit from what it wrote.
748/// Whether a byte is a boundary depends on the bytes around its newline alone,
749/// so a boundary below some limit is the same boundary a scan stopping at that
750/// limit would have found, and the tests hold the two to that.
751fn scan_boundaries(buf: &[u8], out: &mut Vec<usize>) -> bool {
752    out.clear();
753    let mut from = 0;
754    while let Some(rel) = crate::byte_simd::find(&buf[from..], b"\n") {
755        let nl = from + rel;
756        from = nl + 1;
757        if from < buf.len() && !buf[from].is_ascii_whitespace() && !newline_held(buf, nl) {
758            out.push(from);
759        }
760    }
761    buf.last() == Some(&b'\n') && newline_held(buf, buf.len() - 1)
762}
763
764/// The largest safe split boundary below `limit`, read from the boundaries
765/// [`scan_boundaries`] wrote and whether a token may run past the buffer's
766/// last newline.
767///
768/// `end_held` answers for the buffer's end, which only a limit reaching the end
769/// reads - so the two agree wherever it is consulted.
770fn boundary_below(
771    buf: &[u8],
772    bounds: &[usize],
773    end_held: bool,
774    limit: usize,
775    tail: bool,
776) -> Option<usize> {
777    if tail && limit >= buf.len() && !end_held && buf.last() == Some(&b'\n') {
778        return Some(buf.len());
779    }
780    let above = bounds.partition_point(|&b| b < limit);
781    (above > 0).then(|| bounds[above - 1])
782}
783
784/// The largest safe split boundary in `buf` strictly below `limit`: a
785/// significant byte directly after a newline no token spans. At such a point
786/// the prefix lexes exactly as it would in the whole input. With `tail`, and
787/// `limit` reaching the buffer's end, the end is a boundary too when the buffer
788/// ends in a newline no token may span: the only token that could span that cut
789/// is a run of whitespace, which the pattern then never reads, so a line's
790/// match commits as soon as the line ends. Returns `None` when none exists
791/// below `limit`.
792///
793/// Walks for one limit and answers it. Held as the reference the collected
794/// form is checked against, a limit at a time over inputs that quote, escape
795/// and leave a quote open: a stream reads the collected form, so the two
796/// agreeing is what says the collected form cuts where this would.
797#[cfg(test)]
798fn last_safe_boundary(buf: &[u8], limit: usize, tail: bool) -> Option<usize> {
799    let mut best: Option<usize> = None;
800    for i in 1..limit.min(buf.len()) {
801        if buf[i - 1] == b'\n' && !buf[i].is_ascii_whitespace() && !newline_held(buf, i - 1) {
802            best = Some(i);
803        }
804    }
805    if tail && limit >= buf.len() && buf.last() == Some(&b'\n') && !newline_held(buf, buf.len() - 1) {
806        return Some(buf.len());
807    }
808    best
809}
810
811#[cfg(test)]
812mod tests {
813    use super::*;
814    use crate::engine::scan;
815    use crate::parser::parse;
816
817    #[test]
818    fn one_scan_of_the_boundaries_answers_what_a_scan_per_limit_answers() {
819        // `commit_bounded` reads the collected boundaries where it used to scan
820        // the buffer again, and a disagreement between the two is a buffer cut
821        // in the wrong place - which is the failure both reverted attempts on
822        // this path made. Every limit, past the buffer's end as well, and both
823        // tail flags, over inputs that put a newline after a quote its line
824        // does not close, carry a string over a line by a backslash, hold a
825        // newline in a char literal, end on a quote before the last newline,
826        // and escape a quote inside a string.
827        let cases: [&[u8]; 10] = [
828            b"alpha\nbeta\ngamma\n",
829            b"a \"quoted\nline\" b\nc\n",
830            b"\"unclosed\nstill open\n",
831            b"esc \"a\\\"b\nc\" d\ne\n",
832            b"a \"carried \\\nover\" b\nc\n",
833            b"c = '\n' ;\nnext\n",
834            b"ends on a quote '\n",
835            b"\n\n\n  \nx\n",
836            b"no newline at all",
837            b"",
838        ];
839        for buf in cases {
840            let mut bounds = Vec::new();
841            let in_quote = scan_boundaries(buf, &mut bounds);
842            for limit in 0..=buf.len() + 2 {
843                for tail in [false, true] {
844                    assert_eq!(
845                        boundary_below(buf, &bounds, in_quote, limit, tail),
846                        last_safe_boundary(buf, limit, tail),
847                        "{:?} at limit {limit}, tail {tail}",
848                        String::from_utf8_lossy(buf)
849                    );
850                }
851            }
852        }
853    }
854
855    #[test]
856    fn spectral_pattern_streams_equivalently() {
857        // A spectral atom reads a field built over the whole buffer, so
858        // draining a committed prefix moves the reading. The scanner must
859        // defer every commit for such a pattern.
860        let mut input = String::new();
861        for i in 0..200 {
862            input.push_str(&format!("fn f{i}(a,b){{let c=a+b;return c*2;}}\n"));
863            input.push_str("the quick brown fox jumps over the lazy dog again and again\n");
864        }
865        for chunk in [64usize, 512, 4096] {
866            assert_stream_equiv("\\F{texture:code}", &input, chunk);
867        }
868    }
869
870    /// Feed `input` in fixed-size chunks and assert the streaming match
871    /// set equals the whole-input scan, byte for byte.
872    fn assert_stream_equiv(pattern_src: &str, input: &str, chunk: usize) {
873        let pat = parse(pattern_src).expect("pattern parses");
874        let whole = scan(&pat, input.as_bytes());
875        let bytes = input.as_bytes();
876        let chunks: Vec<&[u8]> = bytes.chunks(chunk.max(1)).collect();
877        let streamed = scan_chunked(&pat, chunks);
878        assert_eq!(
879            streamed, whole,
880            "pattern {pattern_src:?} chunk={chunk} differs from whole-input scan"
881        );
882    }
883
884    #[test]
885    fn committed_matches_drain_once_and_finish_returns_the_rest() {
886        let pat = parse("\\N").expect("pattern parses");
887        let whole = b"one 1 two 2\nthree 3 four 4\n";
888        let mut s = StreamScanner::new(pat.clone());
889        s.push(b"one 1 two 2\nthree 3 fo");
890        let early = s.drain_committed();
891        assert!(!early.is_empty(), "the line boundary settles the matches before it");
892        assert!(s.drain_committed().is_empty(), "a drain hands each match over once");
893        s.push(b"ur 4\n");
894        let mut all = early;
895        all.extend(s.drain_committed());
896        all.extend(s.finish());
897        assert_eq!(all, scan(&pat, whole));
898    }
899
900    /// A stream rebased between pushes reports its later matches from the new
901    /// zero, and those plus the distance it moved are the whole input's.
902    #[test]
903    fn a_rebased_stream_finds_the_matches_the_whole_input_holds() {
904        let pat = parse("\\N").expect("pattern parses");
905        let whole = b"one 1 two 2\nthree 3 four 4\nfive 5\n";
906        let mut s = StreamScanner::new(pat.clone());
907        let mut all: Vec<(usize, usize)> = Vec::new();
908        let mut origin = 0usize;
909        for piece in whole.chunks(5) {
910            s.push(piece);
911            all.extend(s.drain_committed().into_iter().map(|m| (origin + m.start(), origin + m.end())));
912            origin += s.rebase();
913        }
914        all.extend(s.finish().into_iter().map(|m| (origin + m.start(), origin + m.end())));
915        let expected: Vec<(usize, usize)> = scan(&pat, whole).iter().map(|m| (m.start(), m.end())).collect();
916        assert_eq!(all, expected);
917    }
918
919    /// Under a declared shape, a stream finds what a whole-input scan under
920    /// the same shape finds, at every chunk size.
921    #[test]
922    fn a_stream_under_declared_shapes_finds_what_the_whole_scan_does() {
923        let mut shapes = crate::custom::ShapeSet::new();
924        shapes.declare("order = `[A-Z]{3}-[0-9]{4}`", crate::custom::Precedence::Before).expect("declare the shape");
925        let pat = crate::parser::parse_with_shapes("\\{order}", &shapes).expect("pattern parses");
926        let input = b"a ABC-1234 b\nXYZ-0007 c DEF-9000\nnone here\nQRS-0001\n";
927        let expected = crate::engine::scan_with_shapes(&pat, input, &shapes);
928        assert!(!expected.is_empty());
929        for chunk in [1usize, 2, 3, 7, 64] {
930            let mut s = StreamScanner::with_shapes(pat.clone(), shapes.clone());
931            let mut got = Vec::new();
932            for piece in input.chunks(chunk) {
933                s.push(piece);
934                got.extend(s.drain_committed());
935            }
936            got.extend(s.finish());
937            assert_eq!(got, expected, "chunk {chunk}");
938        }
939    }
940
941    #[test]
942    fn metamorphic_equivalence_across_patterns_and_chunk_sizes() {
943        let cases: &[(&str, &str)] = &[
944            ("\\N \\W", "weight 12 kg\nlen 5 m\nmass 9 g\n"),
945            ("\\W:x =x", "the the cat\ndog dog ran\nfoo bar baz\n"),
946            ("<\\W:t>.*</=t>", "<a>x</a>\n<b>yy</b>\n<c>z</c>\n"),
947            (". ~\"END\"", "begin here END\nmore lines END now\nlast END\n"),
948            ("\\W\\B(.*)", "call f(g(x))\nrun h(k(y))\ntail\n"),
949            (".*", "anything at all\ngoes here\n"),
950            ("@2 \\W", "a, hello, c\nd, world, f\n"),
951            ("\\I", "10.0.0.1 host\n192.168.1.1 ok\nfe80::1 v6\n"),
952        ];
953        for (pat, input) in cases {
954            for chunk in [1usize, 2, 3, 5, 7, 13, 64, 1000] {
955                assert_stream_equiv(pat, input, chunk);
956            }
957        }
958    }
959
960    /// Strings and char literals at every chunk size stream to the whole
961    /// input's matches, with LF and with CRLF line endings: a quote its line
962    /// does not close, a string carried over a line by a backslash, a char
963    /// literal holding the quote the other kind would open with and one holding
964    /// a newline, and a line ending on a quote.
965    #[test]
966    fn quoted_tokens_stream_as_the_whole_input_reads_them() {
967        let input = concat!(
968            "say \"no close on this line\n",
969            "and \"a string\" here\n",
970            "let s = \"carried \\\nover\" ;\n",
971            "c = '\"' ; d = \"x\"\n",
972            "e = '\n' ; f = \"y\"\n",
973            "ends on a quote '\n",
974            "'last' \"line\"\n",
975        );
976        let crlf = input.replace('\n', "\r\n");
977        for text in [input, crlf.as_str()] {
978            for pattern in ["\\Q", "\\W", "\\Q \\W"] {
979                for chunk in [1usize, 2, 3, 5, 7, 11, 16, 64] {
980                    assert_stream_equiv(pattern, text, chunk);
981                }
982            }
983        }
984    }
985
986    #[test]
987    fn guarded_pattern_with_forward_literal_across_a_seam() {
988        // The guard literal lands in a later chunk than the match start;
989        // a scanner that committed early would wrongly reject the match.
990        let input = "alpha\nbeta\nthe needle is END\n";
991        for chunk in [1usize, 3, 6, 9, 20] {
992            assert_stream_equiv(". ~\"END\"", input, chunk);
993        }
994    }
995
996    #[test]
997    fn a_required_literal_arriving_late_still_finds_its_match() {
998        // While a byte string every match must contain is absent, the scanner
999        // skips the lex and keeps every byte. A match must contain that string
1000        // but need not begin at it, so one can start among the retained bytes
1001        // and reach a literal that arrives chunks later - and a scanner that
1002        // dropped bytes on the literal's own length would lose exactly those.
1003        // Small chunks put the literal's own bytes across a join as well.
1004        let input = "aaa bbb ccc\nddd eee fff\nneedle tail\nggg needle more\n";
1005        for chunk in [1usize, 2, 3, 7, 16, 64] {
1006            assert_stream_equiv("\"needle\" \\W", input, chunk);
1007        }
1008        // And where it never arrives, the stream must agree that there is
1009        // nothing: the skip must not invent a match any more than lose one.
1010        for chunk in [1usize, 3, 16] {
1011            assert_stream_equiv("\"needle\" \\W", "aaa bbb\nccc ddd\neee fff\n", chunk);
1012        }
1013    }
1014
1015    #[test]
1016    fn match_spanning_many_chunks_is_not_split() {
1017        // A single `.*` match runs the length of a line that is far
1018        // longer than the chunk size; it must not be cut into pieces.
1019        let line = "x ".repeat(200);
1020        let input = format!("{line}\n{line}\n");
1021        for chunk in [1usize, 4, 16, 64] {
1022            assert_stream_equiv(".*", &input, chunk);
1023        }
1024    }
1025
1026    #[test]
1027    fn no_newline_input_still_equivalent() {
1028        // No safe boundary exists, so everything retains until finish;
1029        // the result must still match the whole-input scan.
1030        assert_stream_equiv("\\W:x =x", "the the cat dog dog", 3);
1031    }
1032}