use crate::ast::Pattern;
use crate::engine::Span;
use crate::pattern_set::PatternSet;
use crate::token::Token;
enum Source {
One(Pattern),
Set(Box<PatternSet>),
}
pub struct StreamScanner {
source: Source,
buf: Vec<u8>,
base: usize,
out: Vec<(usize, Span)>,
defers_commit: bool,
peak: usize,
tail: bool,
span_limit: Option<usize>,
committed_through: usize,
shapes: crate::custom::ShapeSet,
declared: crate::custom::ShapeSet,
lex_ws: crate::parallel_lex::TokenWorkspace,
bounds: Vec<usize>,
}
impl StreamScanner {
#[must_use]
pub fn new(pattern: Pattern) -> Self {
Self::with_shapes(pattern, crate::custom::ShapeSet::new())
}
#[must_use]
pub fn with_shapes(pattern: Pattern, shapes: crate::custom::ShapeSet) -> Self {
let span_limit = pattern.max_tokens();
let defers_commit = pattern.depends_on_more_than_its_lines() || !settles(&pattern, span_limit);
let tail = !pattern.reads_whitespace();
let mut scanner = Self::over(Source::One(pattern), defers_commit, tail, span_limit);
if !shapes.is_empty() {
if let Source::One(pattern) = &scanner.source {
scanner.shapes = shapes.with_library_shapes(&pattern.library_kinds());
}
scanner.declared = shapes;
}
scanner
}
#[must_use]
pub fn over_set(set: PatternSet) -> Self {
let span_limit = set.max_tokens();
let defers_commit = set.depends_on_more_than_its_lines()
|| !set.patterns().iter().all(|p| settles(p, span_limit));
let tail = !set.reads_whitespace();
Self::over(Source::Set(Box::new(set)), defers_commit, tail, span_limit)
}
fn over(source: Source, defers_commit: bool, tail: bool, span_limit: Option<usize>) -> Self {
let shapes = match &source {
Source::One(pattern) => {
crate::custom::ShapeSet::new().with_library_shapes(&pattern.library_kinds())
}
Source::Set(_) => crate::custom::ShapeSet::new(),
};
Self {
shapes,
declared: crate::custom::ShapeSet::new(),
lex_ws: crate::parallel_lex::TokenWorkspace::default(),
bounds: Vec::new(),
source,
buf: Vec::new(),
base: 0,
out: Vec::new(),
defers_commit,
peak: 0,
tail,
span_limit,
committed_through: 0,
}
}
#[must_use]
pub fn commits_early(&self) -> bool {
!self.defers_commit
}
#[must_use]
pub fn base(&self) -> usize {
self.base
}
#[must_use]
pub fn retained(&self) -> usize {
self.buf.len()
}
pub fn rebase(&mut self) -> usize {
assert!(self.out.is_empty(), "a stream is rebased only once its committed matches are drained");
let by = self.base;
self.base = 0;
self.committed_through = self.committed_through.saturating_sub(by);
by
}
#[must_use]
pub fn peak_retained(&self) -> usize {
self.peak
}
pub fn push(&mut self, chunk: &[u8]) {
let taking = crate::trace::phase("streaming: taking the chunk in");
self.buf.extend_from_slice(chunk);
drop(taking);
self.peak = self.peak.max(self.buf.len());
if self.defers_commit {
return;
}
match self.span_limit {
Some(limit) => self.commit_bounded(limit),
None => self.try_commit(),
}
}
fn scan_retained(&self) -> Vec<(usize, Span)> {
let from = self.committed_through.saturating_sub(self.base);
match &self.source {
Source::One(pattern) => crate::engine::scan_with_shapes_from(pattern, &self.buf, &self.declared, from)
.into_iter()
.map(|s| (0, s))
.collect(),
Source::Set(set) => set.scan_from(&self.buf, from),
}
}
fn lex_retained_into(
buf: &[u8],
shapes: &crate::custom::ShapeSet,
ws: &mut crate::parallel_lex::TokenWorkspace,
) {
if shapes.is_empty() {
crate::lexer::lex_into(buf, &mut ws.toks);
} else {
let blobs = crate::lexer::blob_runs(buf);
ws.toks = crate::lexer::lex_with_shapes(buf, &blobs, shapes, 0);
}
}
fn scan_retained_over(&self, toks: &[Token]) -> Vec<(usize, Span)> {
let from = self.committed_through.saturating_sub(self.base);
match &self.source {
Source::One(pattern) => {
if self.shapes.is_empty()
&& let Some(spans) = crate::engine::routed_spans_at(pattern, &self.buf, from)
{
return spans.into_iter().map(|s| (0, s)).collect();
}
crate::engine::scan_over_tokens_from(pattern, &self.buf, toks, from)
.into_iter()
.map(|s| (0, s))
.collect()
}
Source::Set(set) => set.scan_from_over(&self.buf, from, toks),
}
}
fn commit_bounded(&mut self, limit: usize) {
let deciding = crate::trace::phase("streaming: deciding the prefix");
let end_held = scan_boundaries(&self.buf, &mut self.bounds);
let decided = boundary_below(&self.buf, &self.bounds, end_held, self.buf.len(), self.tail);
drop(deciding);
let Some(decided) = decided else {
return;
};
let lexing = crate::trace::phase("streaming: lexing the retained buffer");
Self::lex_retained_into(&self.buf, &self.shapes, &mut self.lex_ws);
let toks = &self.lex_ws.toks;
drop(lexing);
crate::trace::counted(
"streaming: bytes the buffer holds",
u64::try_from(self.buf.len()).expect("a buffer within the counter's width"),
);
crate::trace::counted(
"streaming: bytes the tokens occupy",
u64::try_from(toks.len() * std::mem::size_of::<Token>())
.expect("a token vector within the counter's width"),
);
let gathering = crate::trace::phase("streaming: gathering the significant tokens");
let sig: Vec<usize> = toks
.iter()
.filter(|t| t.is_significant() && t.end() <= decided)
.map(crate::token::Token::start)
.collect();
let index_at = |byte: usize| sig.partition_point(|&s| s < byte);
let mut after = index_at(self.committed_through.saturating_sub(self.base));
drop(gathering);
let scanning = crate::trace::phase("streaming: scanning the retained buffer");
crate::trace::counted(
"streaming: bytes the scans cover",
u64::try_from(self.buf.len() - self.committed_through.saturating_sub(self.base))
.expect("a buffer within the counter's width"),
);
let found = self.scan_retained_over(toks);
drop(scanning);
for (member, m) in found {
if m.end() > decided || index_at(m.start()) + limit > sig.len() {
break;
}
self.out.push((member, shifted(m, self.base)));
self.committed_through = self.base + m.end();
after = index_at(m.end());
}
let dropping = crate::trace::phase("streaming: dropping the settled bytes");
let retain = after.max((sig.len() + 1).saturating_sub(limit));
let retain_byte = if retain < sig.len() { sig[retain] } else { decided };
let cut = if retain_byte >= self.buf.len() {
self.buf.len()
} else {
boundary_below(&self.buf, &self.bounds, end_held, retain_byte + 1, self.tail)
.unwrap_or(0)
};
crate::trace::counted(
"streaming: bytes the drain moves",
u64::try_from(self.buf.len() - cut).expect("a buffer within the counter's width"),
);
self.buf.drain(..cut);
self.base += cut;
drop(dropping);
}
pub fn drain_committed(&mut self) -> Vec<Span> {
self.drain_committed_with_members().into_iter().map(|(_, s)| s).collect()
}
pub fn drain_committed_with_members(&mut self) -> Vec<(usize, Span)> {
std::mem::take(&mut self.out)
}
#[must_use]
pub fn finish(self) -> Vec<Span> {
self.finish_with_members().into_iter().map(|(_, s)| s).collect()
}
#[must_use]
pub fn finish_with_members(mut self) -> Vec<(usize, Span)> {
for (member, m) in self.scan_retained() {
self.out.push((member, shifted(m, self.base)));
}
self.out
}
fn try_commit(&mut self) {
let lexing = crate::trace::phase("streaming: lexing the retained buffer");
Self::lex_retained_into(&self.buf, &self.shapes, &mut self.lex_ws);
let toks = &self.lex_ws.toks;
drop(lexing);
let Some(settled) = self.settled_through(toks) else {
return;
};
let scanning = crate::trace::phase("streaming: scanning the retained buffer");
crate::trace::counted(
"streaming: bytes the scans cover",
u64::try_from(self.buf.len() - self.committed_through.saturating_sub(self.base))
.expect("a buffer within the counter's width"),
);
let matches = self.scan_retained_over(toks);
drop(scanning);
let _dropping = crate::trace::phase("streaming: dropping the settled bytes");
let end_held = scan_boundaries(&self.buf, &mut self.bounds);
if matches.is_empty() {
if let Some(cut) = boundary_below(&self.buf, &self.bounds, end_held, settled, false) {
self.buf.drain(..cut);
self.base += cut;
}
return;
}
let spans: Vec<Span> = matches.iter().map(|&(_, s)| s).collect();
let Some(cut) = largest_uncrossed_boundary(
&self.buf,
&self.bounds,
end_held,
&spans,
settled,
false,
) else {
return;
};
for &(member, m) in &matches {
if m.end() <= cut {
self.out.push((member, shifted(m, self.base)));
}
}
self.buf.drain(..cut);
self.base += cut;
}
fn settled_through(&self, toks: &[Token]) -> Option<usize> {
let _walking = crate::trace::phase("streaming: the walk over the tokens");
let of = |p: &Pattern| {
crate::nfa::earliest_unsettled(p, &self.buf, toks).map(|at| at.unwrap_or(self.buf.len()))
};
match &self.source {
Source::One(pattern) => of(pattern),
Source::Set(set) => {
let mut least = self.buf.len();
for p in set.patterns() {
least = least.min(of(p)?);
}
Some(least)
}
}
}
}
const REBASE_AT: usize = 1 << 30;
pub struct HeldStream {
scanner: StreamScanner,
held: Vec<u8>,
held_base: usize,
pushed: usize,
origin: usize,
lines_before: Option<usize>,
}
pub struct Ended {
pub matches: Vec<(usize, Span)>,
pub held: Vec<u8>,
pub base: usize,
pub lines_before: Option<usize>,
}
impl HeldStream {
#[must_use]
pub fn new(scanner: StreamScanner, origin: usize, lines_before: Option<usize>) -> Self {
HeldStream { scanner, held: Vec::new(), held_base: 0, pushed: 0, origin, lines_before }
}
#[must_use]
pub fn commits_early(&self) -> bool {
self.scanner.commits_early()
}
#[must_use]
pub fn held(&self) -> &[u8] {
&self.held
}
#[must_use]
pub fn base(&self) -> usize {
self.origin + self.held_base
}
#[must_use]
pub fn lines_before(&self) -> Option<usize> {
self.lines_before
}
#[must_use]
pub fn settled(&self) -> usize {
self.scanner.base() - self.held_base
}
pub fn push(&mut self, bytes: &[u8]) -> Vec<(usize, Span)> {
self.settle();
self.held.extend_from_slice(bytes);
let from = self.pushed - self.held_base;
self.scanner.push(&self.held[from..]);
self.pushed = self.held_base + self.held.len();
let committed = self.scanner.drain_committed_with_members();
local(committed, self.held_base)
}
#[must_use]
pub fn finish(mut self) -> Ended {
self.settle();
let HeldStream { scanner, held, held_base, origin, lines_before, .. } = self;
let matches = local(scanner.finish_with_members(), held_base);
Ended { matches, held, base: origin + held_base, lines_before }
}
fn settle(&mut self) {
let keep_from = self.scanner.base();
if keep_from > self.held_base {
let dropped = keep_from - self.held_base;
if let Some(lines) = self.lines_before.as_mut() {
*lines += crate::byte_simd::count_byte(&self.held[..dropped], b'\n');
}
self.held.drain(..dropped);
self.held_base = keep_from;
}
if self.held_base >= REBASE_AT {
let by = self.scanner.rebase();
self.origin += by;
self.held_base -= by;
self.pushed -= by;
}
}
}
fn local(committed: Vec<(usize, Span)>, held_base: usize) -> Vec<(usize, Span)> {
let at = |offset: usize| {
u32::try_from(offset - held_base).expect("a retained window is narrower than a span's width")
};
committed.into_iter().map(|(member, s)| (member, Span { start: at(s.start()), end: at(s.end()) })).collect()
}
fn settles(pattern: &Pattern, span_limit: Option<usize>) -> bool {
span_limit.is_some() || crate::nfa::earliest_unsettled(pattern, b"", &[]).is_some()
}
fn shifted(span: Span, base: usize) -> Span {
let at = |offset: usize| {
u32::try_from(offset + base)
.expect("a stream's absolute offset fits the span width as an input's does")
};
Span { start: at(span.start()), end: at(span.end()) }
}
#[must_use]
pub fn scan_chunked<'a>(pattern: &Pattern, chunks: impl IntoIterator<Item = &'a [u8]>) -> Vec<Span> {
let mut s = StreamScanner::new(pattern.clone());
for c in chunks {
s.push(c);
}
s.finish()
}
fn largest_uncrossed_boundary(
buf: &[u8],
bounds: &[usize],
end_held: bool,
matches: &[Span],
limit: usize,
tail: bool,
) -> Option<usize> {
let mut cut = boundary_below(buf, bounds, end_held, limit, tail)?;
loop {
if matches.iter().any(|m| m.start() < cut && m.end() > cut) {
cut = boundary_below(buf, bounds, end_held, cut, tail)?;
continue;
}
return Some(cut);
}
}
fn newline_held(buf: &[u8], nl: usize) -> bool {
if nl == 0 {
return false;
}
match buf[nl - 1] {
b'\'' => nl + 1 >= buf.len() || crate::lexer::char_literal_end(buf, nl - 1) == Some(nl + 2),
_ => crate::lexer::backslash_before_newline(buf, nl),
}
}
fn scan_boundaries(buf: &[u8], out: &mut Vec<usize>) -> bool {
out.clear();
let mut from = 0;
while let Some(rel) = crate::byte_simd::find(&buf[from..], b"\n") {
let nl = from + rel;
from = nl + 1;
if from < buf.len() && !buf[from].is_ascii_whitespace() && !newline_held(buf, nl) {
out.push(from);
}
}
buf.last() == Some(&b'\n') && newline_held(buf, buf.len() - 1)
}
fn boundary_below(
buf: &[u8],
bounds: &[usize],
end_held: bool,
limit: usize,
tail: bool,
) -> Option<usize> {
if tail && limit >= buf.len() && !end_held && buf.last() == Some(&b'\n') {
return Some(buf.len());
}
let above = bounds.partition_point(|&b| b < limit);
(above > 0).then(|| bounds[above - 1])
}
#[cfg(test)]
fn last_safe_boundary(buf: &[u8], limit: usize, tail: bool) -> Option<usize> {
let mut best: Option<usize> = None;
for i in 1..limit.min(buf.len()) {
if buf[i - 1] == b'\n' && !buf[i].is_ascii_whitespace() && !newline_held(buf, i - 1) {
best = Some(i);
}
}
if tail && limit >= buf.len() && buf.last() == Some(&b'\n') && !newline_held(buf, buf.len() - 1) {
return Some(buf.len());
}
best
}
#[cfg(test)]
mod tests {
use super::*;
use crate::engine::scan;
use crate::parser::parse;
#[test]
fn one_scan_of_the_boundaries_answers_what_a_scan_per_limit_answers() {
let cases: [&[u8]; 10] = [
b"alpha\nbeta\ngamma\n",
b"a \"quoted\nline\" b\nc\n",
b"\"unclosed\nstill open\n",
b"esc \"a\\\"b\nc\" d\ne\n",
b"a \"carried \\\nover\" b\nc\n",
b"c = '\n' ;\nnext\n",
b"ends on a quote '\n",
b"\n\n\n \nx\n",
b"no newline at all",
b"",
];
for buf in cases {
let mut bounds = Vec::new();
let in_quote = scan_boundaries(buf, &mut bounds);
for limit in 0..=buf.len() + 2 {
for tail in [false, true] {
assert_eq!(
boundary_below(buf, &bounds, in_quote, limit, tail),
last_safe_boundary(buf, limit, tail),
"{:?} at limit {limit}, tail {tail}",
String::from_utf8_lossy(buf)
);
}
}
}
}
#[test]
fn spectral_pattern_streams_equivalently() {
let mut input = String::new();
for i in 0..200 {
input.push_str(&format!("fn f{i}(a,b){{let c=a+b;return c*2;}}\n"));
input.push_str("the quick brown fox jumps over the lazy dog again and again\n");
}
for chunk in [64usize, 512, 4096] {
assert_stream_equiv("\\F{texture:code}", &input, chunk);
}
}
fn assert_stream_equiv(pattern_src: &str, input: &str, chunk: usize) {
let pat = parse(pattern_src).expect("pattern parses");
let whole = scan(&pat, input.as_bytes());
let bytes = input.as_bytes();
let chunks: Vec<&[u8]> = bytes.chunks(chunk.max(1)).collect();
let streamed = scan_chunked(&pat, chunks);
assert_eq!(
streamed, whole,
"pattern {pattern_src:?} chunk={chunk} differs from whole-input scan"
);
}
#[test]
fn committed_matches_drain_once_and_finish_returns_the_rest() {
let pat = parse("\\N").expect("pattern parses");
let whole = b"one 1 two 2\nthree 3 four 4\n";
let mut s = StreamScanner::new(pat.clone());
s.push(b"one 1 two 2\nthree 3 fo");
let early = s.drain_committed();
assert!(!early.is_empty(), "the line boundary settles the matches before it");
assert!(s.drain_committed().is_empty(), "a drain hands each match over once");
s.push(b"ur 4\n");
let mut all = early;
all.extend(s.drain_committed());
all.extend(s.finish());
assert_eq!(all, scan(&pat, whole));
}
#[test]
fn a_rebased_stream_finds_the_matches_the_whole_input_holds() {
let pat = parse("\\N").expect("pattern parses");
let whole = b"one 1 two 2\nthree 3 four 4\nfive 5\n";
let mut s = StreamScanner::new(pat.clone());
let mut all: Vec<(usize, usize)> = Vec::new();
let mut origin = 0usize;
for piece in whole.chunks(5) {
s.push(piece);
all.extend(s.drain_committed().into_iter().map(|m| (origin + m.start(), origin + m.end())));
origin += s.rebase();
}
all.extend(s.finish().into_iter().map(|m| (origin + m.start(), origin + m.end())));
let expected: Vec<(usize, usize)> = scan(&pat, whole).iter().map(|m| (m.start(), m.end())).collect();
assert_eq!(all, expected);
}
#[test]
fn a_stream_under_declared_shapes_finds_what_the_whole_scan_does() {
let mut shapes = crate::custom::ShapeSet::new();
shapes.declare("order = `[A-Z]{3}-[0-9]{4}`", crate::custom::Precedence::Before).expect("declare the shape");
let pat = crate::parser::parse_with_shapes("\\{order}", &shapes).expect("pattern parses");
let input = b"a ABC-1234 b\nXYZ-0007 c DEF-9000\nnone here\nQRS-0001\n";
let expected = crate::engine::scan_with_shapes(&pat, input, &shapes);
assert!(!expected.is_empty());
for chunk in [1usize, 2, 3, 7, 64] {
let mut s = StreamScanner::with_shapes(pat.clone(), shapes.clone());
let mut got = Vec::new();
for piece in input.chunks(chunk) {
s.push(piece);
got.extend(s.drain_committed());
}
got.extend(s.finish());
assert_eq!(got, expected, "chunk {chunk}");
}
}
#[test]
fn metamorphic_equivalence_across_patterns_and_chunk_sizes() {
let cases: &[(&str, &str)] = &[
("\\N \\W", "weight 12 kg\nlen 5 m\nmass 9 g\n"),
("\\W:x =x", "the the cat\ndog dog ran\nfoo bar baz\n"),
("<\\W:t>.*</=t>", "<a>x</a>\n<b>yy</b>\n<c>z</c>\n"),
(". ~\"END\"", "begin here END\nmore lines END now\nlast END\n"),
("\\W\\B(.*)", "call f(g(x))\nrun h(k(y))\ntail\n"),
(".*", "anything at all\ngoes here\n"),
("@2 \\W", "a, hello, c\nd, world, f\n"),
("\\I", "10.0.0.1 host\n192.168.1.1 ok\nfe80::1 v6\n"),
];
for (pat, input) in cases {
for chunk in [1usize, 2, 3, 5, 7, 13, 64, 1000] {
assert_stream_equiv(pat, input, chunk);
}
}
}
#[test]
fn quoted_tokens_stream_as_the_whole_input_reads_them() {
let input = concat!(
"say \"no close on this line\n",
"and \"a string\" here\n",
"let s = \"carried \\\nover\" ;\n",
"c = '\"' ; d = \"x\"\n",
"e = '\n' ; f = \"y\"\n",
"ends on a quote '\n",
"'last' \"line\"\n",
);
let crlf = input.replace('\n', "\r\n");
for text in [input, crlf.as_str()] {
for pattern in ["\\Q", "\\W", "\\Q \\W"] {
for chunk in [1usize, 2, 3, 5, 7, 11, 16, 64] {
assert_stream_equiv(pattern, text, chunk);
}
}
}
}
#[test]
fn guarded_pattern_with_forward_literal_across_a_seam() {
let input = "alpha\nbeta\nthe needle is END\n";
for chunk in [1usize, 3, 6, 9, 20] {
assert_stream_equiv(". ~\"END\"", input, chunk);
}
}
#[test]
fn a_required_literal_arriving_late_still_finds_its_match() {
let input = "aaa bbb ccc\nddd eee fff\nneedle tail\nggg needle more\n";
for chunk in [1usize, 2, 3, 7, 16, 64] {
assert_stream_equiv("\"needle\" \\W", input, chunk);
}
for chunk in [1usize, 3, 16] {
assert_stream_equiv("\"needle\" \\W", "aaa bbb\nccc ddd\neee fff\n", chunk);
}
}
#[test]
fn match_spanning_many_chunks_is_not_split() {
let line = "x ".repeat(200);
let input = format!("{line}\n{line}\n");
for chunk in [1usize, 4, 16, 64] {
assert_stream_equiv(".*", &input, chunk);
}
}
#[test]
fn no_newline_input_still_equivalent() {
assert_stream_equiv("\\W:x =x", "the the cat dog dog", 3);
}
}