use crate::lexer::{
PairedSignificant, Seams, Significant, lex_chunk_into, lex_chunk_paired_into,
lex_chunk_significant_into,
};
use crate::token::{Token, TokenKind};
#[must_use]
pub fn chunk_blobs(blobs: &[(usize, usize)], s: usize, e: usize) -> Vec<(usize, usize)> {
blobs
.iter()
.filter(|&&(a, b)| a >= s && b <= e)
.map(|&(a, b)| (a - s, b - s))
.collect()
}
const PARALLEL_LEX_THRESHOLD: usize = 64 * 1024;
const CHUNKS_A_CORE: usize = 16;
const LEAST_CHUNK_BYTES: usize = 8 * 1024;
#[must_use]
pub fn chunks_for(n: usize, cores: usize) -> usize {
(n / LEAST_CHUNK_BYTES).clamp(cores * 4, cores * CHUNKS_A_CORE)
}
pub const QUOTE_SCAN_MIN_LEAF: usize = 32 * 1024;
pub const QUOTE_SCAN_PARALLEL_BYTES: usize = 1_000_000;
fn parallel_lex_threshold() -> usize {
static THRESHOLD: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
*THRESHOLD.get_or_init(|| match std::env::var("TREX_PARALLEL_LEX_THRESHOLD") {
Err(std::env::VarError::NotPresent) => PARALLEL_LEX_THRESHOLD,
Err(e) => panic!("TREX_PARALLEL_LEX_THRESHOLD is set but unreadable: {e}"),
Ok(v) => v.parse::<usize>().unwrap_or_else(|e| {
panic!("TREX_PARALLEL_LEX_THRESHOLD is set to {v:?}, which is not a byte count: {e}")
}),
})
}
#[must_use]
pub fn active_parallel_threshold() -> usize {
parallel_lex_threshold()
}
#[derive(Default)]
pub struct TokenWorkspace {
parts: Vec<(Vec<Token>, Seams)>,
pub toks: Vec<Token>,
}
thread_local! {
static HELD: std::cell::Cell<TokenWorkspace> = std::cell::Cell::new(TokenWorkspace::default());
}
pub fn lex_parallel_held<R>(input: &[u8], f: impl FnOnce(&[Token]) -> R) -> R {
if input.len() < parallel_lex_threshold() {
return f(&crate::lexer::lex(input));
}
let mut ws = HELD.take();
lex_parallel_into(input, &mut ws);
let out = f(&ws.toks);
HELD.set(ws);
out
}
#[must_use]
pub fn lex_parallel(input: &[u8]) -> Vec<Token> {
let mut ws = TokenWorkspace::default();
lex_parallel_into(input, &mut ws);
ws.toks
}
pub fn lex_parallel_into(input: &[u8], ws: &mut TokenWorkspace) {
use flynnel::JobPlan;
let TokenWorkspace { parts, toks } = ws;
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
let (bounds, blobs) = bounds_and_blobs(
input,
cores,
(Some("the lex: the boundaries"), Some("the lex: the blob table")),
);
if bounds.len() <= 2 {
toks.clear();
toks.reserve(input.len() / crate::lexer::TOKEN_BYTES_ESTIMATE);
lex_chunk_into(input, &blobs, toks, &mut Seams::default());
return;
}
let ranges: Vec<(usize, usize)> = bounds.windows(2).map(|w| (w[0], w[1])).collect();
if parts.len() < ranges.len() {
parts.resize_with(ranges.len(), Default::default);
}
let parts = &mut parts[..ranges.len()];
for (&(s, e), (chunk, seams)) in ranges.iter().zip(parts.iter_mut()) {
chunk.clear();
chunk.reserve((e - s) / crate::lexer::TOKEN_BYTES_ESTIMATE);
seams.open.clear();
seams.close.clear();
}
let avg_chunk_bytes = input.len() / ranges.len().max(1);
let per_chunk_ns = (avg_chunk_bytes as u64 * 6).min(u32::MAX as u64) as u32;
let plan = JobPlan::new(0, ranges.len() as u32)
.with_leaf_shape(flynnel::LeafShape::Streaming)
.with_estimated_per_item_ns(per_chunk_ns);
let lexing = crate::trace::phase("the lex: a chunk a leaf");
over_ranges_longest_first(&plan, &ranges, parts, |k, (chunk, seams)| {
let (s, e) = ranges[k];
lex_chunk_into(&input[s..e], &chunk_blobs(&blobs, s, e), chunk, seams);
});
drop(lexing);
stitch_parallel_into(parts, &ranges, toks);
}
#[derive(Default)]
pub struct PairedWorkspace {
parts: Vec<(PairedSignificant, Seams)>,
}
pub fn lex_paired_parts_held<R>(
input: &[u8],
f: impl FnOnce(&[(PairedSignificant, Seams)]) -> R,
) -> R {
let mut ws = HELD_PAIRED.take();
let n = fill_paired_parts(input, &mut ws);
let out = f(&ws.parts[..n]);
HELD_PAIRED.set(ws);
out
}
fn fill_paired_parts(input: &[u8], ws: &mut PairedWorkspace) -> usize {
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
let (bounds, blobs) =
bounds_and_blobs(input, cores, (Some("the paired lex: the boundaries"), None));
let ranges: Vec<(usize, usize)> = if bounds.len() <= 2 {
vec![(0, input.len())]
} else {
bounds.windows(2).map(|w| (w[0], w[1])).collect()
};
lex_paired_leaves(input, &ranges, &blobs, ws);
join_paired_seams(&mut ws.parts[..ranges.len()]);
ranges.len()
}
fn lex_paired_leaves(
input: &[u8],
ranges: &[(usize, usize)],
blobs: &[(usize, usize)],
ws: &mut PairedWorkspace,
) {
use flynnel::JobPlan;
if ws.parts.len() < ranges.len() {
ws.parts
.resize_with(ranges.len(), || (PairedSignificant::with_base(0, 0), Seams::default()));
}
let _lexing = crate::trace::phase("the paired lex, a mate a token");
let parts = &mut ws.parts[..ranges.len()];
for (&(s, e), (part, seam)) in ranges.iter().zip(parts.iter_mut()) {
part.reset(s);
part.reserve((e - s) / crate::lexer::TOKEN_BYTES_ESTIMATE);
seam.open.clear();
seam.close.clear();
}
if let [(s, e)] = *ranges {
let (part, seam) = &mut parts[0];
lex_chunk_paired_into(&input[s..e], &chunk_blobs(blobs, s, e), part, seam);
return;
}
let lexed: usize = ranges.iter().map(|&(s, e)| e - s).sum();
let avg_chunk_bytes = lexed / ranges.len().max(1);
let per_chunk_ns = (avg_chunk_bytes as u64 * 6).min(u32::MAX as u64) as u32;
let plan = JobPlan::new(0, ranges.len() as u32)
.with_leaf_shape(flynnel::LeafShape::Streaming)
.with_estimated_per_item_ns(per_chunk_ns);
over_ranges_longest_first(&plan, ranges, &mut *parts, |k, (part, seam)| {
let (s, e) = ranges[k];
lex_chunk_paired_into(&input[s..e], &chunk_blobs(blobs, s, e), part, seam);
});
}
fn join_paired_seams(parts: &mut [(PairedSignificant, Seams)]) {
let _joining = crate::trace::phase("the seam join");
let mut bases: Vec<u32> = Vec::with_capacity(parts.len());
let mut acc: u32 = 0;
for (part, _) in parts.iter() {
bases.push(acc);
acc = acc
.checked_add(u32::try_from(part.mates.len()).expect("a token index within the stored width"))
.expect("a token index within the stored width");
}
let shifting = crate::trace::phase("the seam join: every mate shifted by its base");
for ((part, _), &base) in parts.iter_mut().zip(bases.iter()) {
if base == 0 {
continue;
}
for m in &mut part.mates {
if *m != crate::lexer::NO_MATE {
*m += base;
}
}
}
drop(shifting);
let mut events: Vec<(u32, usize, usize, bool, crate::token::BracketKind)> = Vec::new();
for (ci, (part, seam)) in parts.iter().enumerate() {
let base = bases[ci];
let (mut oi, mut cj) = (0usize, 0usize);
while oi < seam.open.len() || cj < seam.close.len() {
let take_open =
seam.open.get(oi).is_some_and(|&o| cj >= seam.close.len() || o < seam.close[cj]);
let local = if take_open { seam.open[oi] } else { seam.close[cj] };
if let Some((is_open, bk)) = TokenKind::bracket_of_code(part.parts.kinds[local]) {
events.push((base + local as u32, ci, local, is_open, bk));
}
if take_open {
oi += 1;
} else {
cj += 1;
}
}
}
let mut open_stack: Vec<(u32, usize, usize, crate::token::BracketKind)> = Vec::new();
let mut paired: Vec<(usize, usize, usize, usize, u32, u32)> = Vec::new();
for (at, ci, local, is_open, bk) in events {
if is_open {
open_stack.push((at, ci, local, bk));
continue;
}
if let Some(&(open_at, open_ci, open_local, open_bk)) = open_stack.last()
&& open_bk == bk
{
open_stack.pop();
paired.push((open_ci, open_local, ci, local, at, open_at));
}
}
for (open_ci, open_local, close_ci, close_local, close_at, open_at) in paired {
parts[open_ci].0.mates[open_local] = close_at;
parts[close_ci].0.mates[close_local] = open_at;
}
}
thread_local! {
static HELD_PAIRED: std::cell::Cell<PairedWorkspace> =
std::cell::Cell::new(PairedWorkspace::default());
}
#[derive(Default)]
pub struct SignificantWorkspace {
parts: Vec<Significant>,
pub kinds: Vec<u32>,
pub spans: Vec<(u32, u32)>,
}
#[must_use]
pub fn lex_significant_parallel(input: &[u8]) -> (Vec<u32>, Vec<(u32, u32)>) {
let mut ws = SignificantWorkspace::default();
lex_significant_parallel_into(input, &mut ws);
(ws.kinds, ws.spans)
}
thread_local! {
static HELD_SIGNIFICANT: std::cell::Cell<SignificantWorkspace> =
std::cell::Cell::new(SignificantWorkspace::default());
}
pub fn lex_significant_parallel_held<R>(
input: &[u8],
f: impl FnOnce(&[u32], &[(u32, u32)]) -> R,
) -> R {
let mut ws = HELD_SIGNIFICANT.take();
lex_significant_parallel_into(input, &mut ws);
let out = f(&ws.kinds, &ws.spans);
HELD_SIGNIFICANT.set(ws);
out
}
pub fn lex_significant_parts_held<R>(input: &[u8], f: impl FnOnce(&[Significant]) -> R) -> R {
let mut ws = HELD_SIGNIFICANT.take();
let n = fill_significant_parts(input, &mut ws);
let out = f(&ws.parts[..n]);
HELD_SIGNIFICANT.set(ws);
out
}
pub struct OwnedParts {
ws: Option<SignificantWorkspace>,
n: usize,
}
impl OwnedParts {
#[must_use]
pub fn parts(&self) -> &[Significant] {
let ws = self.ws.as_ref().expect("the workspace is taken only as this is dropped");
&ws.parts[..self.n]
}
}
impl Drop for OwnedParts {
fn drop(&mut self) {
if let Some(ws) = self.ws.take() {
HELD_SIGNIFICANT.set(ws);
}
}
}
#[must_use]
pub fn lex_significant_parts_owned(input: &[u8]) -> OwnedParts {
let mut ws = HELD_SIGNIFICANT.take();
let n = fill_significant_parts(input, &mut ws);
OwnedParts { ws: Some(ws), n }
}
fn fill_significant_parts(input: &[u8], ws: &mut SignificantWorkspace) -> usize {
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
let (bounds, blobs) = bounds_and_blobs(
input,
cores,
(Some("the lex in parts: the boundaries"), Some("the lex in parts: the blob table")),
);
if bounds.len() <= 2 {
if ws.parts.is_empty() {
ws.parts.push(Significant::with_base(0, 0));
}
let part = &mut ws.parts[0];
part.reset(0);
part.reserve(input.len() / crate::lexer::TOKEN_BYTES_ESTIMATE);
lex_chunk_significant_into(input, &chunk_blobs(&blobs, 0, input.len()), part);
return 1;
}
let ranges: Vec<(usize, usize)> = bounds.windows(2).map(|w| (w[0], w[1])).collect();
let _leafing = crate::trace::phase("the lex in parts: a chunk a leaf");
lex_significant_leaves(input, &ranges, &blobs, &mut ws.parts).len()
}
pub fn lex_significant_parallel_into(input: &[u8], ws: &mut SignificantWorkspace) {
ws.kinds.clear();
ws.spans.clear();
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
let (bounds, blobs) = bounds_and_blobs(input, cores, (None, None));
if bounds.len() <= 2 {
lex_significant_serial_into(input, 0, input.len(), &blobs, ws);
return;
}
let ranges: Vec<(usize, usize)> = bounds.windows(2).map(|w| (w[0], w[1])).collect();
lex_significant_chunks_into(input, &ranges, &blobs, ws);
}
pub fn lex_significant_range_into(
input: &[u8],
s: usize,
e: usize,
blobs: &[(usize, usize)],
ws: &mut SignificantWorkspace,
) {
ws.kinds.clear();
ws.spans.clear();
lex_significant_range_onto(input, s, e, blobs, ws);
}
pub fn lex_significant_range_onto(
input: &[u8],
s: usize,
e: usize,
blobs: &[(usize, usize)],
ws: &mut SignificantWorkspace,
) {
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
let bounds = if e - s < parallel_lex_threshold() || cores <= 1 {
Vec::new()
} else {
safe_boundaries(&input[s..e], chunks_for(e - s, cores))
};
if bounds.len() <= 2 {
lex_significant_serial_into(input, s, e, blobs, ws);
return;
}
let ranges: Vec<(usize, usize)> = bounds.windows(2).map(|w| (s + w[0], s + w[1])).collect();
lex_significant_chunks_into(input, &ranges, blobs, ws);
}
fn lex_significant_serial_into(
input: &[u8],
s: usize,
e: usize,
blobs: &[(usize, usize)],
ws: &mut SignificantWorkspace,
) {
let SignificantWorkspace { parts, kinds, spans } = ws;
if parts.is_empty() {
parts.push(Significant::with_base(0, 0));
}
let part = &mut parts[0];
part.reset(s);
part.reserve((e - s) / crate::lexer::TOKEN_BYTES_ESTIMATE);
lex_chunk_significant_into(&input[s..e], &chunk_blobs(blobs, s, e), part);
kinds.extend_from_slice(&part.kinds);
spans.extend_from_slice(&part.spans);
}
fn longest_first(ranges: &[(usize, usize)]) -> Vec<usize> {
let mut order: Vec<usize> = (0..ranges.len()).collect();
order.sort_by_key(|&k| std::cmp::Reverse(ranges[k].1 - ranges[k].0));
order
}
#[allow(unsafe_code)]
fn over_ranges_longest_first<T: Send>(
plan: &flynnel::JobPlan,
ranges: &[(usize, usize)],
parts: &mut [T],
lex: impl Fn(usize, &mut T) + Sync,
) {
use flynnel::sched::par_iter::for_each_chunk_indexed_min_leaf;
assert_eq!(parts.len(), ranges.len(), "a part a range");
let mut order = longest_first(ranges);
let parts_at = parts.as_mut_ptr() as usize;
for_each_chunk_indexed_min_leaf(plan, &mut order, 1, |_, slots| {
for &k in slots.iter() {
let part = unsafe { &mut *(parts_at as *mut T).add(k) };
lex(k, part);
}
});
}
#[allow(unsafe_code)]
fn lex_significant_leaves<'a>(
input: &[u8],
ranges: &[(usize, usize)],
blobs: &[(usize, usize)],
parts: &'a mut Vec<Significant>,
) -> &'a mut [Significant] {
use flynnel::JobPlan;
use flynnel::sched::par_iter::for_each_chunk_indexed_min_leaf;
use std::sync::atomic::{AtomicU64, Ordering};
if parts.len() < ranges.len() {
parts.resize_with(ranges.len(), || Significant::with_base(0, 0));
}
let parts = &mut parts[..ranges.len()];
for (&(s, e), part) in ranges.iter().zip(parts.iter_mut()) {
part.reset(s);
part.reserve((e - s) / crate::lexer::TOKEN_BYTES_ESTIMATE);
}
let lexed: usize = ranges.iter().map(|&(s, e)| e - s).sum();
let avg_chunk_bytes = lexed / ranges.len().max(1);
let per_chunk_ns = (avg_chunk_bytes as u64 * 6).min(u32::MAX as u64) as u32;
let plan = JobPlan::new(0, ranges.len() as u32)
.with_leaf_shape(flynnel::LeafShape::Streaming)
.with_estimated_per_item_ns(per_chunk_ns);
let dispatched = crate::trace::keeping().then(std::time::Instant::now);
let leaves = AtomicU64::new(0);
let leaf_nanos = AtomicU64::new(0);
let longest = AtomicU64::new(0);
let latest_start = AtomicU64::new(0);
let mut order = longest_first(ranges);
let parts_at = parts.as_mut_ptr() as usize;
for_each_chunk_indexed_min_leaf(&plan, &mut order, 1, |_, slots| {
let clock = dispatched.map(|dispatched| (dispatched.elapsed(), std::time::Instant::now()));
for &k in slots.iter() {
let (s, e) = ranges[k];
let slot = unsafe { &mut *(parts_at as *mut Significant).add(k) };
lex_chunk_significant_into(&input[s..e], &chunk_blobs(blobs, s, e), slot);
}
if let Some((since_dispatch, began)) = clock {
let took = nanos(began.elapsed());
leaves.fetch_add(1, Ordering::Relaxed);
leaf_nanos.fetch_add(took, Ordering::Relaxed);
longest.fetch_max(took, Ordering::Relaxed);
latest_start.fetch_max(nanos(since_dispatch), Ordering::Relaxed);
}
});
if dispatched.is_some() {
crate::trace::counted("the lex in parts: leaves", leaves.into_inner());
crate::trace::counted("the lex in parts: leaf nanoseconds, summed", leaf_nanos.into_inner());
crate::trace::counted("the lex in parts: the longest leaf, nanoseconds", longest.into_inner());
crate::trace::counted(
"the lex in parts: the latest leaf start, nanoseconds",
latest_start.into_inner(),
);
}
parts
}
fn nanos(d: std::time::Duration) -> u64 {
u64::try_from(d.as_nanos()).expect("a leaf shorter than the five hundred years the counter holds")
}
#[allow(unsafe_code, clippy::uninit_vec)]
fn lex_significant_chunks_into(
input: &[u8],
ranges: &[(usize, usize)],
blobs: &[(usize, usize)],
ws: &mut SignificantWorkspace,
) {
use flynnel::JobPlan;
use flynnel::sched::par_iter::for_each_chunk_indexed_min_leaf;
let SignificantWorkspace { parts, kinds, spans } = ws;
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
let parts = lex_significant_leaves(input, ranges, blobs, parts);
let total: usize = parts.iter().map(|p| p.kinds.len()).sum();
if total == 0 {
return;
}
let base = kinds.len();
let mut offsets = Vec::with_capacity(parts.len());
let mut acc = base;
for p in parts.iter() {
offsets.push(acc);
acc += p.kinds.len();
}
kinds.reserve(total);
spans.reserve(total);
unsafe {
kinds.set_len(base + total);
spans.set_len(base + total);
}
let kinds_addr = kinds.as_mut_ptr() as usize;
let spans_addr = spans.as_mut_ptr() as usize;
let parts: &[Significant] = parts;
let mut fan: Vec<u8> = vec![0; parts.len()];
let min_leaf = parts.len().div_ceil(cores * 4).max(1);
let plan = JobPlan::new(0, parts.len() as u32).with_leaf_shape(flynnel::LeafShape::Streaming);
for_each_chunk_indexed_min_leaf(&plan, &mut fan, min_leaf, |start, slots| {
let dst_kinds = kinds_addr as *mut u32;
let dst_spans = spans_addr as *mut (u32, u32);
for k in 0..slots.len() {
let ci = start + k;
let off = offsets[ci];
let part = &parts[ci];
unsafe {
std::ptr::copy_nonoverlapping(part.kinds.as_ptr(), dst_kinds.add(off), part.kinds.len());
std::ptr::copy_nonoverlapping(part.spans.as_ptr(), dst_spans.add(off), part.spans.len());
}
}
});
}
#[must_use]
pub fn stitch_parallel(parts: &[(Vec<Token>, Seams)], ranges: &[(usize, usize)]) -> Vec<Token> {
let mut out = Vec::new();
stitch_parallel_into(parts, ranges, &mut out);
out
}
#[allow(unsafe_code, clippy::uninit_vec)]
pub fn stitch_parallel_into(
parts: &[(Vec<Token>, Seams)],
ranges: &[(usize, usize)],
out: &mut Vec<Token>,
) {
use flynnel::JobPlan;
use flynnel::sched::par_iter::for_each_chunk_indexed_min_leaf;
let _copying = crate::trace::phase("the stitch into one array");
out.clear();
let nchunks = parts.len();
let total: usize = parts.iter().map(|(toks, _)| toks.len()).sum();
if total == 0 {
return;
}
let mut offsets = Vec::with_capacity(nchunks);
let mut acc = 0usize;
for (toks, _) in parts {
offsets.push(acc);
acc += toks.len();
}
out.reserve(total);
unsafe {
out.set_len(total);
}
let out_addr = out.as_mut_ptr() as usize;
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
let mut fan: Vec<u8> = vec![0; nchunks];
let min_leaf = nchunks.div_ceil(cores * 4).max(1);
let plan = JobPlan::new(0, nchunks as u32).with_leaf_shape(flynnel::LeafShape::Streaming);
for_each_chunk_indexed_min_leaf(&plan, &mut fan, min_leaf, |start, slots| {
let dst = out_addr as *mut Token;
for k in 0..slots.len() {
let ci = start + k;
let base = ranges[ci].0;
let off = offsets[ci];
for (j, t) in parts[ci].0.iter().enumerate() {
let mut tok = *t;
tok.shift(base);
tok.set_mate(tok.mate().map(|m| off + m));
unsafe {
dst.add(off + j).write(tok);
}
}
}
});
let mut open_stack: Vec<usize> = Vec::new();
for (ci, (_, seam)) in parts.iter().enumerate() {
let off = offsets[ci];
let (mut oi, mut cj) = (0, 0);
while oi < seam.open.len() || cj < seam.close.len() {
let open_next = seam.open.get(oi).is_some_and(|&o| cj >= seam.close.len() || o < seam.close[cj]);
if open_next {
open_stack.push(off + seam.open[oi]);
oi += 1;
} else {
let idx = off + seam.close[cj];
if let TokenKind::Close(bk) = out[idx].kind
&& let Some(&open_idx) = open_stack.last()
&& out[open_idx].kind == TokenKind::Open(bk)
{
open_stack.pop();
out[open_idx].set_mate(Some(idx));
out[idx].set_mate(Some(open_idx));
}
cj += 1;
}
}
}
}
fn bounds_and_blobs(
input: &[u8],
cores: usize,
phases: (Option<&'static str>, Option<&'static str>),
) -> (Vec<usize>, Vec<(usize, usize)>) {
if input.len() < parallel_lex_threshold() || cores <= 1 {
return (Vec::new(), crate::lexer::blob_runs(input));
}
let (bounding, tabling) = phases;
let bounds = {
let _timed = bounding.map(crate::trace::phase);
safe_boundaries(input, chunks_for(input.len(), cores))
};
let blobs = {
let _timed = tabling.map(crate::trace::phase);
if bounds.len() <= 2 {
crate::lexer::blob_runs(input)
} else {
crate::lexer::blob_runs_parallel(input)
}
};
(bounds, blobs)
}
#[must_use]
pub fn safe_boundaries(input: &[u8], target_chunks: usize) -> Vec<usize> {
let n = input.len();
let quotes = quoted_spans_across(input, QUOTE_SCAN_MIN_LEAF);
let holding = |p: usize| {
let k = quotes.partition_point(|&(open, _)| open <= p);
(k > 0 && quotes[k - 1].1 >= p).then(|| quotes[k - 1])
};
let mut picked: Vec<usize> = Vec::new();
let (mut at_newline, mut at_space, mut uncut, mut runs_walked) = (0u64, 0u64, 0u64, 0u64);
for k in 1..target_chunks {
let mut from = k * n / target_chunks;
let limit = ((k + 1) * n / target_chunks).min(n);
let before = picked.len();
while from < limit {
let Some(rel) = crate::byte_simd::find(&input[from..limit], b"\n") else {
break;
};
let newline = from + rel;
let held = |at: usize| crate::lexer::char_literal_end(input, at) == Some(newline + 2);
if (newline >= 1 && held(newline - 1)) || (newline >= 2 && held(newline - 2)) {
from = newline + 1;
continue;
}
let mut b = newline + 1;
while b < n && input[b].is_ascii_whitespace() {
b += 1;
}
if b >= n {
break;
}
match holding(newline).or_else(|| holding(b)) {
Some((_, close)) => from = close + 1,
None => {
picked.push(b);
at_newline += 1;
break;
}
}
}
if picked.len() == before {
let mut from = k * n / target_chunks;
while from < limit {
let w = from + crate::byte_simd::nonspace_run(&input[from..limit]);
if w >= limit {
break;
}
runs_walked += 1;
let p = w + crate::byte_simd::space_run(&input[w..]);
if let Some((_, close)) = holding(w).or_else(|| holding(p)) {
from = close + 1;
continue;
}
let unit_after_a_number = p == w + 1
&& input[w] == b' '
&& input[w - 1].is_ascii_digit()
&& crate::quantity::spaced_symbol_len(input, p).is_some();
let safe = p < n
&& w > 0
&& !matches!(input[w - 1], b'\'' | b'\\')
&& !input[p].is_ascii_digit()
&& !matches!(input[p], b'-' | b'+')
&& !unit_after_a_number;
if safe {
picked.push(p);
at_space += 1;
break;
}
from = p;
}
if picked.len() == before {
uncut += 1;
}
}
}
crate::trace::counted("the boundaries: divisions cut at a newline", at_newline);
crate::trace::counted("the boundaries: divisions cut at a whitespace run", at_space);
crate::trace::counted("the boundaries: divisions left uncut", uncut);
crate::trace::counted("the boundaries: whitespace runs the fallback walked", runs_walked);
picked.sort_unstable();
picked.dedup();
let mut bounds = Vec::with_capacity(picked.len() + 2);
bounds.push(0);
bounds.extend(picked);
bounds.push(n);
bounds.dedup();
bounds
}
#[must_use]
pub fn quoted_spans(input: &[u8]) -> Vec<(usize, usize)> {
match quoted_spans_where(input, |_| Some(true)) {
Some(spans) => spans,
None => unreachable!("a rule that settles every quote abandons no reading"),
}
}
pub(crate) fn quoted_spans_where(
input: &[u8],
mut opens: impl FnMut(usize) -> Option<bool>,
) -> Option<Vec<(usize, usize)>> {
let read = quote_scan_where(input, 0, input.len(), &mut opens)?;
Some(read.spans)
}
#[derive(Clone, Default)]
struct QuoteRead {
spans: Vec<(usize, usize)>,
from: usize,
}
fn quote_scan_where(
input: &[u8],
from: usize,
to: usize,
opens: &mut impl FnMut(usize) -> Option<bool>,
) -> Option<QuoteRead> {
let mut spans = Vec::new();
let mut quotes = crate::byte_simd::QuotePositions::new(&input[..to]);
let mut closes = crate::byte_simd::CloseOrNewlinePositions::new(input);
let mut from = from;
while let Some(p) = quotes.next_at_or_after(from) {
from = p + 1;
if input[p] == b'\'' {
if let Some(end) = crate::lexer::char_literal_end(input, p) {
if opens(p)? {
from = end;
}
} else if let Some(end) = crate::lexer::single_quoted_end(input, p)
&& opens(p)?
{
spans.push((p, end - 1));
from = end;
}
} else if let Some(end) = crate::lexer::double_quoted_end_at(input, p, &mut closes)
&& opens(p)?
{
spans.push((p, end - 1));
from = end;
}
}
Some(QuoteRead { spans, from })
}
fn quote_scan(input: &[u8], from: usize, to: usize) -> QuoteRead {
match quote_scan_where(input, from, to, &mut |_| Some(true)) {
Some(read) => read,
None => unreachable!("a rule that settles every quote abandons no reading"),
}
}
fn quote_piece_starts(input: &[u8], leaf_bytes: usize) -> Option<Vec<usize>> {
let n = input.len();
let reach = leaf_bytes.min(QUOTE_SCAN_MIN_LEAF);
let mut starts = vec![0usize];
let mut at = leaf_bytes;
while at < n {
let window_end = (at + reach).min(n);
let mut from = at;
let mut start = None;
while let Some(rel) = crate::byte_simd::find(&input[from..window_end], b"\n") {
let newline = from + rel;
if !crate::lexer::backslash_before_newline(input, newline) {
start = Some(newline + 1);
break;
}
from = newline + 1;
}
match start {
Some(s) if s < n => {
starts.push(s);
at = s + leaf_bytes;
}
Some(_) => break,
None if window_end == n => break,
None => return None,
}
}
Some(starts)
}
#[must_use]
pub fn quoted_spans_across(input: &[u8], least_a_leaf: usize) -> Vec<(usize, usize)> {
if input.len() < QUOTE_SCAN_PARALLEL_BYTES {
return quoted_spans(input);
}
quoted_spans_in_pieces(input, quote_piece_bytes(input.len(), least_a_leaf))
}
#[must_use]
pub fn quote_scan_pieces(input: &[u8], least_a_leaf: usize) -> usize {
if input.len() < QUOTE_SCAN_PARALLEL_BYTES {
return 1;
}
quote_piece_starts(input, quote_piece_bytes(input.len(), least_a_leaf)).map_or(1, |s| s.len())
}
fn quote_piece_bytes(n: usize, least_a_leaf: usize) -> usize {
let cores = std::thread::available_parallelism().map_or(1, std::num::NonZero::get);
n.div_ceil(cores * 4).max(least_a_leaf).max(1)
}
fn quoted_spans_in_pieces(input: &[u8], leaf_bytes: usize) -> Vec<(usize, usize)> {
use flynnel::JobPlan;
use flynnel::sched::par_iter::for_each_chunk_indexed_min_leaf;
let n = input.len();
let Some(starts) = quote_piece_starts(input, leaf_bytes) else {
return quoted_spans(input);
};
if starts.len() <= 1 {
return quoted_spans(input);
}
let ends: Vec<usize> = starts[1..].iter().copied().chain(std::iter::once(n)).collect();
let mut read: Vec<QuoteRead> = vec![QuoteRead::default(); starts.len()];
let per_piece_ns = (leaf_bytes as u64 * 150 / 1000).min(u64::from(u32::MAX)) as u32;
let pieces = u32::try_from(starts.len()).expect("a piece count of at most four a core, and one");
let plan = JobPlan::new(0, pieces)
.with_leaf_shape(flynnel::LeafShape::Streaming)
.with_estimated_per_item_ns(per_piece_ns);
for_each_chunk_indexed_min_leaf(&plan, &mut read, 1, |base, slots| {
for (i, slot) in slots.iter_mut().enumerate() {
let k = base + i;
*slot = quote_scan(input, starts[k], ends[k]);
}
});
let mut spans = Vec::new();
let mut from = 0usize;
for (k, piece) in read.iter().enumerate() {
let entered_past_start;
let piece = if from > starts[k] {
entered_past_start = quote_scan(input, from, ends[k]);
&entered_past_start
} else {
piece
};
spans.extend_from_slice(&piece.spans);
from = piece.from;
}
spans
}
pub(crate) struct QuoteScan<'a> {
input: &'a [u8],
spans: Vec<(usize, usize)>,
next_double: Option<usize>,
next_single: Option<usize>,
searched_to: usize,
settled_to: usize,
closes: crate::byte_simd::CloseOrNewlinePositions<'a>,
unsettled: bool,
done: bool,
}
impl<'a> QuoteScan<'a> {
pub(crate) fn new(input: &'a [u8]) -> Self {
QuoteScan {
input,
spans: Vec::new(),
next_double: None,
next_single: None,
searched_to: 0,
settled_to: 0,
closes: crate::byte_simd::CloseOrNewlinePositions::new(input),
unsettled: false,
done: false,
}
}
pub(crate) fn unsettled(&self) -> bool {
self.unsettled
}
pub(crate) fn spans(&self) -> &[(usize, usize)] {
&self.spans
}
#[cfg(test)]
pub(crate) fn searched_to(&self) -> usize {
self.searched_to
}
pub(crate) fn ensure(&mut self, upto: usize, opens: &mut impl FnMut(usize) -> Option<bool>) {
let n = self.input.len();
let input = self.input;
let next = |from: usize, to: usize, quote: u8| {
(from < to)
.then(|| crate::byte_simd::find(&input[from..to], &[quote]).map(|r| from + r))
.flatten()
};
while !self.done && self.settled_to <= upto {
let horizon = upto.saturating_add(1).max(self.searched_to.saturating_mul(2)).min(n);
if self.next_double.is_none() && self.next_single.is_none() && self.searched_to < horizon
{
self.next_double = next(self.searched_to, horizon, b'"');
self.next_single = next(self.searched_to, horizon, b'\'');
self.searched_to = horizon;
}
let p = match (self.next_double, self.next_single) {
(Some(d), Some(s)) => d.min(s),
(Some(q), None) | (None, Some(q)) => q,
(None, None) => {
if horizon == n {
self.done = true;
}
self.settled_to = horizon;
break;
}
};
if p > upto {
self.settled_to = p;
break;
}
let mut from = p + 1;
let read = if input[p] == b'\'' {
crate::lexer::char_literal_end(input, p)
.map(|end| (end, false))
.or_else(|| crate::lexer::single_quoted_end(input, p).map(|end| (end, true)))
} else {
crate::lexer::double_quoted_end_at(input, p, &mut self.closes).map(|end| (end, true))
};
if let Some((end, string)) = read {
match opens(p) {
Some(true) => {
if string {
self.spans.push((p, end - 1));
}
from = end;
}
Some(false) => {}
None => {
self.unsettled = true;
self.done = true;
return;
}
}
}
self.settled_to = from;
if from > self.searched_to {
self.searched_to = from;
}
if self.next_double.is_some_and(|q| q < from) {
self.next_double = next(from, self.searched_to, b'"');
}
if self.next_single.is_some_and(|q| q < from) {
self.next_single = next(from, self.searched_to, b'\'');
}
}
}
}
#[must_use]
pub fn quoted_spans_resumed(input: &[u8], step: usize) -> Vec<(usize, usize)> {
let mut scan = QuoteScan::new(input);
let mut always = |_: usize| Some(true);
for upto in (0..input.len()).step_by(step.max(1)) {
scan.ensure(upto, &mut always);
}
scan.ensure(input.len(), &mut always);
scan.spans().to_vec()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::lexer::{lex, lex_chunk};
#[test]
fn owning_the_parts_gives_the_lent_stream_and_hands_the_workspace_back() {
fn flat(parts: &[Significant]) -> Vec<(u32, u32, u32)> {
parts
.iter()
.flat_map(|p| p.kinds.iter().zip(&p.spans).map(|(k, s)| (*k, s.0, s.1)))
.collect()
}
let mut big = String::new();
for i in 0..20_000u32 {
big.push_str(&format!("let value_{i} = {} ; call_{i}(a, \"q\") ;\n", i * 7));
}
for input in [b"" as &[u8], b"x", b"a 1 b 2\n", big.as_bytes()] {
let lent = lex_significant_parts_held(input, flat);
let owned = lex_significant_parts_owned(input);
assert_eq!(flat(owned.parts()), lent, "{} bytes", input.len());
drop(owned);
assert_eq!(lex_significant_parts_held(input, flat), lent);
}
let outer = lex_significant_parts_owned(b"a 1 b 2\n");
let inner = lex_significant_parts_owned(b"a 1 b 2\n");
assert_eq!(flat(outer.parts()), flat(inner.parts()));
}
#[test]
fn the_quote_scan_in_pieces_is_the_whole_scan() {
let own: &[u8] = include_bytes!("parallel_lex.rs");
let lexer_src: &[u8] = include_bytes!("lexer.rs");
let mut edges = String::new();
for i in 0..40 {
edges.push_str(&format!("let s{i} = \"first line\\\nsecond \\\" line\\\nthird \\\\\";\n"));
edges.push_str("let c = '\n';\n");
edges.push_str("say 'held \\\nover' done\n");
edges.push_str(&format!("let r{i} = \"crlf\\\r\nnext\";\r\n"));
edges.push_str("say 'crlf \\\r\nheld' done\r\n");
edges.push_str("macro \\\nnext\n");
edges.push_str("don't \"quote runs\n on here\n");
}
let mut open_end = edges.clone();
open_end.push_str("tail \"never closed\nat all\n");
for (name, input, leaves) in [
("this file", own, &[256usize, 1000, 4096, 1 << 20][..]),
("lexer.rs", lexer_src, &[256, 1000, 4096, 1 << 20]),
("the edges", edges.as_bytes(), &[48, 64, 256, 1000]),
("a quote left open at the end", open_end.as_bytes(), &[48, 64, 256]),
("no newline", b"a \"b\" 'c' \"d" as &[u8], &[1, 4, 64]),
("empty", b"", &[1, 64]),
] {
let whole = quoted_spans(input);
for &leaf in leaves {
assert_eq!(quoted_spans_in_pieces(input, leaf), whole, "{name}, {leaf} bytes a piece");
}
assert_eq!(quoted_spans_across(input, 1), whole, "{name}, across the cores");
}
let input = edges.as_bytes();
let whole = quoted_spans(input);
let mut cut_after = std::collections::BTreeSet::new();
for leaf in 40..=104 {
assert_eq!(quoted_spans_in_pieces(input, leaf), whole, "the edges, {leaf} bytes a piece");
let starts = quote_piece_starts(input, leaf).expect("every line of the edges is shorter than a piece");
for &s in &starts[1..] {
assert_eq!(input[s - 1], b'\n', "a piece at {s} starts after a newline");
assert!(
!crate::lexer::backslash_before_newline(input, s - 1),
"a piece at {s} starts after an escaped newline"
);
cut_after.insert(s - 1);
}
}
let held = input.windows(3).position(|w| w == b"'\n'").expect("the edges hold a char literal holding a newline");
assert!(cut_after.contains(&(held + 1)), "no cut fell after the newline a char literal holds");
}
#[test]
fn newlines_too_far_apart_to_cut_at_are_read_whole() {
let mut sparse = "x".repeat(5000);
sparse.push_str(" \"a\" 'b'\n");
sparse.push_str(&"y".repeat(5000));
let input = sparse.as_bytes();
assert_eq!(quote_piece_starts(input, 1000), None, "a piece's search came up empty and was not refused");
assert_eq!(quoted_spans_in_pieces(input, 1000), quoted_spans(input));
let lines: String = (0..500).map(|i| format!("line {i} \"q\" 'c' ok\n")).collect();
let cut = quote_piece_starts(lines.as_bytes(), 1000).expect("newlines every line are within a piece");
assert!(cut.len() >= 5, "a dense input divides into pieces: {} starts", cut.len());
}
#[test]
fn the_resumable_quote_scan_agrees_with_the_eager_one() {
let mut quoted = String::new();
for i in 0..400 {
match i % 5 {
0 => quoted.push_str(&format!("let a_{i} = \"str {i}\" ;\n")),
1 => quoted.push_str(&format!("call_{i}('c', \"esc \\\" still in\", {i}) ;\n")),
2 => quoted.push_str(&format!("url_{i} = 'http://x/{i}' ;\n")),
3 => quoted.push_str(&format!("plain_{i} = {i} ;\n")),
_ => quoted.push_str(&format!("q_{i} = \"trailing \\\\\" ;\n")),
}
}
let mut bare = String::new();
for i in 0..400 {
bare.push_str(&format!("let value_{i} = {i} ; call_{i}(alpha, beta) ;\n"));
}
let mut late = bare.clone();
late.push_str("tail = \"only string\" ;\n");
for (name, src) in [("quoted", "ed), ("quote-free", &bare), ("quote-late", &late)] {
let input = src.as_bytes();
let eager = quoted_spans(input);
assert_eq!(
eager.is_empty(),
name == "quote-free",
"{name}: the corpus must exercise what it is for"
);
let mut scan = QuoteScan::new(input);
let mut always = |_: usize| Some(true);
for upto in [0usize, 1, 40, 200, 1000, 5000, input.len() / 2, input.len() - 1, input.len()]
{
scan.ensure(upto, &mut always);
assert!(!scan.unsettled(), "{name}: this corpus settles every quote");
let got = scan.spans();
assert!(got.len() <= eager.len(), "{name}: decided more spans than exist at {upto}");
assert_eq!(
got,
&eager[..got.len()],
"{name}: the decided prefix at {upto} is not the eager scan's"
);
let owed = eager.iter().filter(|s| s.0 <= upto).count();
assert!(
got.len() >= owed,
"{name} at {upto}: decided {} of the {owed} spans owed",
got.len()
);
}
scan.ensure(input.len(), &mut always);
assert_eq!(scan.spans(), &eager[..], "{name}: the finished scans differ");
}
}
#[test]
fn the_resumable_quote_scan_reads_no_further_than_it_was_asked() {
let mut src = String::new();
for i in 0..4000 {
src.push_str(&format!("let value_{i} = {i} ; call_{i}(alpha, beta) ;\n"));
}
let input = src.as_bytes();
assert!(input.len() > 100_000, "the input is long enough for the ask to be a small part");
let mut scan = QuoteScan::new(input);
let mut always = |_: usize| Some(true);
scan.ensure(64, &mut always);
assert!(
scan.searched_to() <= 65,
"an ask about byte 64 read {} bytes of {}",
scan.searched_to(),
input.len()
);
assert!(scan.spans().is_empty());
scan.ensure(65, &mut always);
assert!(
scan.searched_to() >= 130,
"a second ask read only to {}, so ascending asks do not double",
scan.searched_to()
);
}
#[test]
fn a_quote_past_the_ask_is_not_reached_however_far_the_search_ran() {
let mut src = String::new();
for i in 0..4000 {
src.push_str(&format!("let value_{i} = {i} ;\n"));
}
let tail = src.len();
src.push_str("x = 'http://unsettled' ;\n");
let input = src.as_bytes();
let mut scan = QuoteScan::new(input);
let mut rule = |q: usize| (q < tail).then_some(true);
for upto in [0usize, 64, 1000, tail / 2, tail - 1] {
scan.ensure(upto, &mut rule);
assert!(!scan.unsettled(), "an ask about {upto} reached the quote at {tail}");
}
scan.ensure(input.len(), &mut rule);
assert!(scan.unsettled(), "the quote is reached once it is asked about");
}
fn candidates_by_walk(input: &[u8]) -> Vec<usize> {
let mut out = Vec::new();
let mut i = 0;
while i < input.len() {
let b = input[i];
if let Some(end) = crate::lexer::double_quoted_end(input, i) {
i = end;
continue;
}
if i > 0 && input[i - 1].is_ascii_whitespace() && !b.is_ascii_whitespace() {
let mut w = i;
while w > 0 && input[w - 1].is_ascii_whitespace() {
w -= 1;
}
let after_newline = input[i - 1] == b'\n';
let by_whitespace = w > 0
&& !matches!(input[w - 1], b'\'' | b'\\')
&& !b.is_ascii_digit()
&& !matches!(b, b'-' | b'+');
if after_newline || by_whitespace {
out.push(i);
}
}
i += 1;
}
out
}
#[test]
fn single_quoted_strings_lex_the_same_chunked_and_whole() {
let corpus = concat!(
"let a = 'one two' ;\n",
"let b = 'has \"double\" inside' ; x\n",
"let c = \"has 'single' inside\" ;\n",
"fn f<'a>(x: &'a str) -> &'a str { x }\n",
"don't stop, it's the '90s ;\n",
"let d = 'never closed\n",
"rock 'n' roll and '\\'' ;\n",
"y = '\\u{10000}\\u{10001}' 'z' ;\n",
"plain line at the end\n",
);
let input = corpus.as_bytes();
assert_identical(input);
let spans = quoted_spans(input);
let quoted: Vec<(usize, usize)> = lex(input)
.iter()
.filter(|t| t.kind == crate::token::TokenKind::Quoted && t.end() - t.start() > 4)
.map(|t| (t.start(), t.end() - 1))
.collect();
for span in "ed {
assert!(spans.contains(span), "the lexer's string {span:?} is not among the reader's {spans:?}");
}
}
#[test]
fn boundaries_fall_only_where_the_walk_allows() {
let corpus = concat!(
"let a = \"one\\\"two\" ;\n",
"let b = \"spans\nlines\nhere\" ; x\n",
"let c = \"ends with backslashes\\\\\" ;\n",
"let d = \"odd\\\\\\\" still open\n",
"not closed yet ;\n",
"\" closed now\n",
"\"opens the line\" ok\n",
"let e = \"carried \\\nover\" ;\n",
"let f = \"carried \\\r\nover crlf\" ;\r\n",
"plain line one\n",
"plain line two\n",
"y = 'q' \"tail without close\n",
"last line\n",
);
let input = corpus.as_bytes();
let allowed = candidates_by_walk(input);
for chunks in [2usize, 3, 5, 8, 13, 40] {
let bounds = safe_boundaries(input, chunks);
assert_eq!(bounds[0], 0);
assert_eq!(*bounds.last().expect("a last bound"), input.len());
for &b in &bounds[1..bounds.len() - 1] {
assert!(allowed.contains(&b), "boundary {b} at {chunks} chunks is not after a newline outside a string");
}
let ranges: Vec<(usize, usize)> = bounds.windows(2).map(|w| (w[0], w[1])).collect();
let blobs = crate::lexer::blob_runs(input);
let parts: Vec<(Vec<Token>, Seams)> = ranges
.iter()
.map(|&(s, e)| lex_chunk(&input[s..e], &chunk_blobs(&blobs, s, e), 16))
.collect();
assert_eq!(stitch_parallel(&parts, &ranges), lex(input), "{chunks} chunks");
}
let one_string = b"\"a\\\\\" b\nc\n";
assert_eq!(quoted_spans(one_string), vec![(0, 4)]);
assert_eq!(quoted_spans(b"\"\\\"\nx"), Vec::new());
assert_eq!(quoted_spans(b"\"a\\\nb\" c"), vec![(0, 5)]);
assert_eq!(quoted_spans(b"no quotes\n"), Vec::new());
}
#[test]
fn a_string_read_past_the_search_hides_the_quotes_inside_it() {
let mut src = String::new();
for i in 0..200 {
src.push_str(&format!("let value_{i} = {i} ;\n"));
}
let open = src.len() + 4;
src.push_str("x = \"a 'held inside' and 'c' still in the string\" ;\n");
for i in 0..200 {
src.push_str(&format!("let later_{i} = 'z' ;\n"));
}
let input = src.as_bytes();
assert_eq!(input[open], b'"');
let eager = quoted_spans(input);
let mut scan = QuoteScan::new(input);
let mut always = |_: usize| Some(true);
scan.ensure(open, &mut always);
assert!(scan.searched_to() > open, "the string was read past the first search");
scan.ensure(input.len(), &mut always);
assert_eq!(scan.spans(), &eager[..], "a quote inside the string was read as an opener");
}
#[test]
fn unicode_text_stitches_byte_identical() {
let line = "caf\u{E9} Gr\u{FC}\u{DF}e \u{03B8} na\u{EF}ve\u{2014}done \u{A0} ok\n";
let big = line.repeat(64);
assert_eq!(lex_chunked_at_every_boundary(big.as_bytes()), lex(big.as_bytes()));
}
fn lex_chunked_at_every_boundary(input: &[u8]) -> Vec<Token> {
let bounds = safe_boundaries(input, input.len().max(1));
if bounds.len() <= 2 {
return lex(input);
}
let ranges: Vec<(usize, usize)> = bounds.windows(2).map(|w| (w[0], w[1])).collect();
let blobs = crate::lexer::blob_runs(input);
let parts: Vec<(Vec<Token>, Seams)> = ranges
.iter()
.map(|&(s, e)| lex_chunk(&input[s..e], &chunk_blobs(&blobs, s, e), 0))
.collect();
stitch_parallel(&parts, &ranges)
}
#[test]
fn whitespace_boundaries_keep_every_recognizer_whole() {
let corpus = concat!(
"card 4111 1111 1111 1111 paid 4111-1111-1111-1111 also 12345678901234567 sum ",
"call +44 20 7123 4567 or +1 555-123-4567 or 212-555-1234 now ",
"at 37.7749, -122.4194 and 55.7558, 37.6173 or 12.5,34.7 here ",
"c = ' ' d = '\\ ' e = 'x' f = \"a string with spaces\" g = \"esc \\\" quote\" ",
"size 10MB and 1.5GiB took 1500ms or 3h20m for $1,234.56 at 50% ",
"mass 5 kg at 3.2 GHz and 40 % or -40\u{b0}C over 5 m/s with 5 items and 3 in a row ",
"on 2024-01-02 at 12:30:45 or 2024-01-02T12:30:00 see /usr/bin and C:\\Users\\x ",
"mail user@example.com http://example.com/a?b=1 v1.2.3 #fff 00:1a:2b:3c:4d:5e ",
"hash d41d8cd98f00b204e9800998ecf8427e blob SGVsbG8gV29ybGQhIQ== end",
);
let input = corpus.as_bytes();
assert!(!input.contains(&b'\n'));
let bounds = safe_boundaries(input, input.len());
assert!(bounds.len() > 20, "the whitespace rule should split a newline-free input: {bounds:?}");
assert_eq!(lex_chunked_at_every_boundary(input), lex(input));
}
#[test]
fn seam_brackets_pair_as_the_serial_lexer_pairs_them() {
for input in [
b"( a\n] b\n) c\n".as_slice(),
b"[ x\n( y\n] z\n) w\n",
b"a )\n( b\n) c\n",
b"{\n[\n(\n)\n]\n}\n",
b"( ]\n)\n",
b") (\n) ]\n",
b"[ ( x\n) y ]\n} z\n",
b"( ( )\n) )\n( \n",
] {
assert_eq!(
lex_chunked_at_every_boundary(input),
lex(input),
"chunked pairing differs from serial on {:?}",
String::from_utf8_lossy(input)
);
}
}
#[test]
fn significant_lex_equals_the_serial_streams_significant_part() {
let grow = |make: &dyn Fn(usize) -> Vec<u8>| {
let mut units = 4000;
loop {
let candidate = make(units);
if candidate.len() > PARALLEL_LEX_THRESHOLD {
break candidate;
}
units *= 2;
}
};
let prose = grow(&prose_and_blobs);
let brackets = grow(&|n| {
let mut v = Vec::new();
for i in 0..n {
v.extend_from_slice(format!("a (b [c {i}] d) e\n(f\ng) {{h}}\n").as_bytes());
}
v
});
let small = b"call +1 555-123-4567 now 4111 1111 1111 1111 pin 37.7749,-122.4194 x".to_vec();
for (name, input) in [("prose and blobs", &prose), ("brackets across lines", &brackets), ("small", &small)]
{
let want = crate::gpu::significant_stream(&lex(input));
let got = lex_significant_parallel(input);
assert!(!want.0.is_empty(), "{name}: nothing to compare");
assert_eq!(got.0, want.0, "{name}: kind codes differ from the serial lexer's");
assert_eq!(got.1, want.1, "{name}: spans differ from the serial lexer's");
}
}
#[test]
fn the_leaves_take_the_longest_range_first_and_write_each_part_once() {
assert_eq!(longest_first(&[(0, 10), (10, 50), (50, 55), (55, 95)]), vec![1, 3, 0, 2]);
assert_eq!(longest_first(&[(0, 4), (4, 8), (8, 12)]), vec![0, 1, 2]);
assert_eq!(longest_first(&[]), Vec::<usize>::new());
let input: &[u8] = include_bytes!("parallel_lex.rs");
let n = input.len();
let cuts = [0, n / 20, n / 20 + 2000, n / 3, n / 3 + 100, n / 2, 3 * n / 4, n];
let ranges: Vec<(usize, usize)> = cuts.windows(2).map(|w| (w[0], w[1])).collect();
let blobs = crate::lexer::blob_runs(input);
let mut parts = Vec::new();
let got = lex_significant_leaves(input, &ranges, &blobs, &mut parts);
assert_eq!(got.len(), ranges.len());
for (k, &(s, e)) in ranges.iter().enumerate() {
let mut want = Significant::with_base(0, 0);
want.reset(s);
lex_chunk_significant_into(&input[s..e], &chunk_blobs(&blobs, s, e), &mut want);
assert!(!want.kinds.is_empty(), "range {k} holds tokens");
assert_eq!(got[k].kinds, want.kinds, "range {k}");
assert_eq!(got[k].spans, want.spans, "range {k}");
}
}
#[test]
fn significant_chunks_at_every_boundary_join_to_the_serial_stream() {
use crate::lexer::lex_chunk_significant;
let input: &[u8] = b"call +1 555-123-4567 now 4111 1111 1111 1111 pin 37.7749,-122.4194 x 'a b' \"q r\" (p [q] r) 10.0.0.1 v1.2.3 end";
let bounds = safe_boundaries(input, input.len());
assert!(bounds.len() > 2, "the walk must split this input: {bounds:?}");
let blobs = crate::lexer::blob_runs(input);
let (mut kinds, mut spans) = (Vec::new(), Vec::new());
for w in bounds.windows(2) {
let (s, e) = (w[0], w[1]);
let part = lex_chunk_significant(&input[s..e], &chunk_blobs(&blobs, s, e), s, 0);
kinds.extend_from_slice(&part.kinds);
spans.extend_from_slice(&part.spans);
}
assert_eq!((kinds, spans), crate::gpu::significant_stream(&lex(input)));
}
#[test]
fn significant_ranges_between_safe_boundaries_join_to_the_whole_lex() {
let mut units = 4000;
let prose = loop {
let candidate = prose_and_blobs(units);
if candidate.len() > 8 * PARALLEL_LEX_THRESHOLD {
break candidate;
}
units *= 2;
};
let whole = lex_significant_parallel(&prose);
let blobs = crate::lexer::blob_runs_parallel(&prose);
let mut ws = SignificantWorkspace::default();
for partitions in [2, 3, 8, 64] {
let bounds = safe_boundaries(&prose, partitions);
assert!(bounds.len() > 2, "{partitions} partitions: the walk must split the input");
let (mut kinds, mut spans) = (Vec::new(), Vec::new());
for w in bounds.windows(2) {
lex_significant_range_into(&prose, w[0], w[1], &blobs, &mut ws);
kinds.extend_from_slice(&ws.kinds);
spans.extend_from_slice(&ws.spans);
}
assert_eq!(kinds, whole.0, "{partitions} partitions: kind codes differ from the whole lex");
assert_eq!(spans, whole.1, "{partitions} partitions: spans differ from the whole lex");
}
}
#[test]
fn held_token_buffers_carry_nothing_between_inputs() {
let mut units = 4000;
let big = loop {
let candidate = prose_and_blobs(units);
if candidate.len() > PARALLEL_LEX_THRESHOLD {
break candidate;
}
units *= 2;
};
let small = b"a (b [c] d) e\n".to_vec();
let mut ws = TokenWorkspace::default();
lex_parallel_into(&big, &mut ws);
assert_eq!(ws.toks, lex(&big));
lex_parallel_into(&small, &mut ws);
assert_eq!(ws.toks, lex(&small));
lex_parallel_into(&big, &mut ws);
assert_eq!(ws.toks, lex(&big));
assert_eq!(lex_parallel_held(&small, <[Token]>::to_vec), lex(&small));
assert_eq!(lex_parallel_held(&big, <[Token]>::len), lex(&big).len());
}
#[test]
fn held_significant_buffers_carry_nothing_between_inputs() {
let mut units = 4000;
let big = loop {
let candidate = prose_and_blobs(units);
if candidate.len() > PARALLEL_LEX_THRESHOLD {
break candidate;
}
units *= 2;
};
let small = b"a (b [c] d) e\n".to_vec();
let mut ws = SignificantWorkspace::default();
for input in [&big, &small, &big] {
lex_significant_parallel_into(input, &mut ws);
let want = crate::gpu::significant_stream(&lex(input));
assert_eq!((ws.kinds.clone(), ws.spans.clone()), want);
}
}
#[test]
fn a_held_lex_inside_a_held_lex_gets_buffers_of_its_own() {
let outer = b"a (b) c\n".repeat(20_000);
let inner = b"x [y] z\n".repeat(10_000);
assert!(
outer.len() > PARALLEL_LEX_THRESHOLD && inner.len() > PARALLEL_LEX_THRESHOLD,
"both lexes must take the pool route, where the buffers are held"
);
let (outer_toks, inner_toks) =
lex_parallel_held(&outer, |o| (o.to_vec(), lex_parallel_held(&inner, <[Token]>::to_vec)));
assert_eq!(outer_toks, lex(&outer));
assert_eq!(inner_toks, lex(&inner));
let (outer_sig, inner_sig) = lex_significant_parallel_held(&outer, |ok, os| {
let inner_sig = lex_significant_parallel_held(&inner, |ik, is| (ik.to_vec(), is.to_vec()));
((ok.to_vec(), os.to_vec()), inner_sig)
});
assert_eq!(outer_sig, crate::gpu::significant_stream(&lex(&outer)));
assert_eq!(inner_sig, crate::gpu::significant_stream(&lex(&inner)));
}
fn prose_and_blobs(lines: usize) -> Vec<u8> {
const ALPHA: &[u8] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
let mut out = Vec::new();
let mut x = 0x1234_5678u64;
for i in 0..lines {
out.extend_from_slice(
format!("the quick brown fox jumps over the lazy dog line {i}\n").as_bytes(),
);
for _ in 0..64 {
x = x
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1_442_695_040_888_963_407);
out.push(ALPHA[((x >> 58) % 64) as usize]);
}
out.push(b'\n');
}
out
}
#[test]
fn blob_bearing_input_stitches_byte_identical() {
assert_identical(&prose_and_blobs(200));
}
#[test]
fn lex_parallel_matches_serial_over_the_threshold() {
let mut units = 4000;
let input = loop {
let candidate = prose_and_blobs(units);
if candidate.len() > PARALLEL_LEX_THRESHOLD {
break candidate;
}
units *= 2;
};
assert_eq!(lex_parallel(&input), lex(&input));
}
fn assert_identical(input: &[u8]) {
assert_eq!(
lex_chunked_at_every_boundary(input),
lex(input),
"chunked lex differs from serial on {:?}",
String::from_utf8_lossy(input)
);
}
#[test]
fn matches_serial_on_line_oriented_input() {
assert_identical(b"alpha beta\n12 34\nword (group)\n");
}
#[test]
fn a_newline_held_in_a_char_literal_is_never_a_cut() {
assert_identical(b"see [[Hummingbird]].''\n'''Hummer''' is a [[marque]] of vehicles\n");
assert_identical(b"let c = '\\\n' ;\nnext line\n");
assert_identical(b"a '\n' b\n'\n'\nc\n");
}
#[test]
fn reconciles_brackets_across_a_chunk_seam() {
let input = b"open (\nmiddle line here\n) close\ntail\n";
assert_identical(input);
let toks = lex_chunked_at_every_boundary(input);
let open = toks.iter().position(|t| matches!(t.kind, TokenKind::Open(_))).unwrap();
let close = toks.iter().position(|t| matches!(t.kind, TokenKind::Close(_))).unwrap();
assert_eq!(toks[open].mate(), Some(close));
assert_eq!(toks[close].mate(), Some(open));
}
#[test]
fn parallel_stitch_matches_serial() {
for input in [
b"alpha beta\n12 34\nword (group)\n".as_slice(),
b"open (\nmiddle line here\n) close\ntail\n",
b"a (b (c) d) e\nf (g) h\n(i\nj)\nk\n",
b"no brackets here\njust words and 123 numbers\n",
b"[x]\n{y}\n(z)\nmix [a (b) c]\n",
] {
let bounds = safe_boundaries(input, input.len().max(1));
if bounds.len() <= 2 {
continue;
}
let ranges: Vec<(usize, usize)> = bounds.windows(2).map(|w| (w[0], w[1])).collect();
let parts: Vec<(Vec<Token>, Seams)> =
ranges.iter().map(|&(s, e)| lex_chunk(&input[s..e], &[], 0)).collect();
assert_eq!(
stitch_parallel(&parts, &ranges),
lex(input),
"parallel stitch differs from serial lex on {:?}",
String::from_utf8_lossy(input)
);
}
}
#[test]
fn never_splits_inside_a_quoted_string_with_newlines() {
let input = b"before\n\"a quoted\nstring with\nnewlines\"\nafter\n";
assert_identical(input);
}
#[test]
fn handles_blank_lines_and_leading_whitespace() {
let input = b"a\n\n\n b\nc\n";
assert_identical(input);
}
#[test]
fn handles_no_trailing_newline() {
assert_identical(b"x\ny\nz");
}
#[test]
fn handles_typed_tokens_per_line() {
let input =
b"2026-06-16 192.168.0.1 a@b.com\nhttp://x.com/p 12:30:45\nfe80::1 plain word\n";
assert_identical(input);
}
#[test]
fn empty_and_tiny_inputs() {
assert_identical(b"");
assert_identical(b"a");
assert_identical(b"\n");
}
#[test]
fn the_paired_parts_pair_what_the_stitched_tokens_pair() {
let mut src = String::from("( {\n");
for i in 0..40_000 {
src.push_str(&format!("f_{i} [ a_{i} ] b_{i} ( c_{i} )\n"));
}
src.push_str("} )\n");
let input = src.as_bytes();
let toks = lex_parallel(input);
let mut sig_of = vec![usize::MAX; toks.len()];
let mut tok_of: Vec<usize> = Vec::new();
for (i, t) in toks.iter().enumerate() {
if t.kind != TokenKind::Whitespace {
sig_of[i] = tok_of.len();
tok_of.push(i);
}
}
lex_paired_parts_held(input, |parts| {
assert!(parts.len() > 1, "the input is long enough to be split");
let mut mates: Vec<u32> = Vec::new();
for (part, _) in parts {
mates.extend_from_slice(&part.mates);
}
assert_eq!(
mates.len(),
tok_of.len(),
"one mate slot per significant token"
);
let mut crossed = 0usize;
for (s, &t) in tok_of.iter().enumerate() {
let want = toks[t].mate().map(|m| sig_of[m]);
let got = (mates[s] != crate::lexer::NO_MATE).then(|| mates[s] as usize);
assert_eq!(got, want, "significant {s} (token {t})");
if let Some(m) = want
&& m.abs_diff(s) > 4
{
crossed += 1;
}
}
assert!(crossed >= 4, "the outer brackets pair across chunks");
});
}
}