use std::ops::Range;
use scrive_core::{
Document, HighlightEngine, Revision, SegmentBoundary, SegmentStart, SegmentTokens, Snapshot,
};
pub(crate) const PARALLEL_MIN_BYTES: u32 = 2 * 1_048_576;
const SEGMENT_MAX_ROWS: u32 = 32_768;
const SPECULATION_BACKOFF: u32 = 128;
const POLL_VERIFY_BUDGET: usize = 4;
const POLL_RERUN_BUDGET: usize = 1;
#[cfg(test)]
thread_local! {
static POLL_RERUNS: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
}
struct Job {
idx: usize,
snapshot: std::sync::Arc<Snapshot>,
rows: Range<u32>,
start: SegmentStart,
spans_for: Range<u32>,
against: Option<SegmentTokens>,
rev: Revision,
}
struct Done {
idx: usize,
seg: SegmentTokens,
rev: Revision,
}
pub(crate) struct HighlightPool {
engine: HighlightEngine,
fresh: SegmentBoundary,
queue: std::sync::Arc<(std::sync::Mutex<std::collections::VecDeque<Job>>, std::sync::Condvar)>,
cur_rev: std::sync::Arc<std::sync::atomic::AtomicU64>,
done_rx: std::sync::mpsc::Receiver<Done>,
_workers: Vec<std::thread::JoinHandle<()>>,
snapshot: std::sync::Arc<Snapshot>,
pub(crate) rev: Revision,
seg_rows: Vec<Range<u32>>,
results: Vec<Option<SegmentTokens>>,
next_verify: usize,
prev_end: Option<SegmentBoundary>,
window: Range<u32>,
pub(crate) active: bool,
}
impl HighlightPool {
pub(crate) fn new(doc: &Document, viewport: Range<u32>) -> Option<Self> {
let engine = doc.highlight_engine()?;
let count = std::thread::available_parallelism()
.map(|n| n.get().saturating_sub(1).max(1))
.unwrap_or(1);
let (done_tx, done_rx) = std::sync::mpsc::channel::<Done>();
let queue: std::sync::Arc<(std::sync::Mutex<std::collections::VecDeque<Job>>, std::sync::Condvar)> =
std::sync::Arc::new((std::sync::Mutex::new(std::collections::VecDeque::new()), std::sync::Condvar::new()));
let cur_rev = std::sync::Arc::new(std::sync::atomic::AtomicU64::new(doc.revision().0));
let workers = (0..count)
.map(|_| {
let engine = engine.clone();
let queue = queue.clone();
let cur_rev = cur_rev.clone();
let done_tx = done_tx.clone();
std::thread::spawn(move || loop {
let job = {
let (lock, cvar) = &*queue;
let mut q = lock.lock().unwrap();
while q.is_empty() {
q = cvar.wait(q).unwrap();
}
q.pop_front().unwrap()
};
if job.rev.0 != cur_rev.load(std::sync::atomic::Ordering::Relaxed) {
continue;
}
let seg = scrive_core::tokenize_segment(
&engine, &job.snapshot, job.rows, job.start, Some(job.spans_for),
job.against.as_ref(),
);
let _ = done_tx.send(Done { idx: job.idx, seg, rev: job.rev });
})
})
.collect();
let fresh = engine.fresh_boundary();
let mut pool = Self {
engine,
fresh,
queue,
cur_rev,
done_rx,
_workers: workers,
snapshot: std::sync::Arc::new(doc.snapshot()),
rev: doc.revision(),
seg_rows: Vec::new(),
results: Vec::new(),
next_verify: 0,
prev_end: None,
window: 0..0,
active: false,
};
pool.start(doc, viewport);
Some(pool)
}
fn start(&mut self, doc: &Document, viewport: Range<u32>) {
self.rev = doc.revision();
self.cur_rev.store(self.rev.0, std::sync::atomic::Ordering::Relaxed);
self.queue.0.lock().unwrap().clear();
self.snapshot = std::sync::Arc::new(doc.snapshot());
let n = self.snapshot.line_count();
self.set_window(viewport, n);
let workers = self._workers.len() as u32;
let segs = (n.div_ceil(SEGMENT_MAX_ROWS)).max(workers).max(1);
let seg_len = n.div_ceil(segs).max(1);
self.seg_rows = (0..n)
.step_by(seg_len as usize)
.map(|s| s..(s + seg_len).min(n))
.collect();
self.results = (0..self.seg_rows.len()).map(|_| None).collect();
self.next_verify = 0;
self.prev_end = None;
self.active = !self.seg_rows.is_empty();
let (lock, cvar) = &*self.queue;
let mut q = lock.lock().unwrap();
for (idx, rows) in self.seg_rows.iter().enumerate() {
q.push_back(Job {
idx,
snapshot: self.snapshot.clone(),
rows: rows.clone(),
start: SegmentStart::Fresh,
spans_for: self.window.clone(),
against: None,
rev: self.rev,
});
}
cvar.notify_all();
}
fn set_window(&mut self, viewport: Range<u32>, n: u32) {
self.window = scrive_core::padded_highlight_window(viewport, n);
}
pub(crate) fn restart(&mut self, doc: &mut Document, viewport: Range<u32>) {
self.start(doc, viewport.clone());
self.speculate(doc, viewport);
}
pub(crate) fn reaim(&mut self, doc: &mut Document, viewport: Range<u32>) {
let n = doc.buffer().line_count();
self.set_window(viewport.clone(), n);
self.speculate(doc, viewport);
}
pub(crate) fn speculate(&self, doc: &mut Document, viewport: Range<u32>) {
if !self.active || self.rev != doc.revision() {
return;
}
let n = self.snapshot.line_count();
let _ = viewport;
let rows = self.window.start.saturating_sub(SPECULATION_BACKOFF)..self.window.end.min(n);
if rows.start >= rows.end {
return;
}
let seg = scrive_core::tokenize_segment(
&self.engine,
&self.snapshot,
rows,
SegmentStart::Fresh,
Some(self.window.clone()),
None,
);
doc.absorb_highlight(self.rev, seg, false);
}
pub(crate) fn poll(&mut self, doc: &mut Document) {
while let Ok(done) = self.done_rx.try_recv() {
if done.rev == self.rev {
self.results[done.idx] = Some(done.seg);
} }
let mut verified = 0usize;
let mut reruns = 0usize;
while self.next_verify < self.results.len() {
if verified >= POLL_VERIFY_BUDGET {
break; }
let Some(seg_ref) = self.results[self.next_verify].as_ref() else { break };
let true_start_fresh =
self.next_verify == 0 || self.prev_end.as_ref() == Some(&self.fresh);
let needs_rerun = seg_ref.started_fresh() && !true_start_fresh;
if needs_rerun && reruns >= POLL_RERUN_BUDGET {
break; }
let seg = self.results[self.next_verify].take().expect("peeked Some");
if needs_rerun {
let start = SegmentStart::After(self.prev_end.clone().expect("prev segment is verified"));
let fixed = scrive_core::tokenize_segment(
&self.engine,
&self.snapshot,
self.seg_rows[self.next_verify].clone(),
start,
None, Some(&seg),
);
let end = fixed.end_boundary().clone();
doc.absorb_highlight(self.rev, fixed, true);
self.prev_end = Some(end);
reruns += 1;
#[cfg(test)]
POLL_RERUNS.with(|c| c.set(c.get() + 1));
} else {
let end = seg.end_boundary().clone();
doc.absorb_highlight(self.rev, seg, true);
self.prev_end = Some(end);
}
self.next_verify += 1;
verified += 1;
}
if self.next_verify >= self.results.len() && self.active {
self.active = false;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use scrive_core::SyntaxDef;
const GRAMMAR: &str = "%YAML 1.2\n---\nscope: source.t\ncontexts:\n main:\n - match: '\"'\n push: string\n string:\n - match: '\"'\n pop: true\n";
fn grammar_doc(source: &str) -> Document {
let mut doc = Document::new(source).expect("fits the u32 offset space");
doc.set_syntax(
SyntaxDef::from_sublime_syntax(GRAMMAR).expect("grammar parses"),
crate::scrive_dark_theme(),
);
doc
}
#[test]
fn poll_verifies_at_most_budget_per_frame() {
const N: usize = 12;
let mut lines = [""; N];
lines[2] = "\"";
lines[4] = "\"";
let source = lines.join("\n");
let mut doc = grammar_doc(&source);
assert_eq!(doc.buffer().line_count(), N as u32, "one row per segment");
let mut pool = HighlightPool::new(&doc, 0..N as u32).expect("grammar → engine");
pool.rev = scrive_core::Revision(u64::MAX);
pool.cur_rev.store(u64::MAX, std::sync::atomic::Ordering::Relaxed);
pool.queue.0.lock().unwrap().clear();
while pool.done_rx.try_recv().is_ok() {}
pool.snapshot = std::sync::Arc::new(doc.snapshot());
pool.seg_rows = (0..N as u32).map(|r| r..r + 1).collect();
pool.window = pool.snapshot.line_count()..pool.snapshot.line_count();
pool.results = pool
.seg_rows
.iter()
.map(|rows| {
Some(scrive_core::tokenize_segment(
&pool.engine,
&pool.snapshot,
rows.clone(),
SegmentStart::Fresh,
None,
None,
))
})
.collect();
pool.next_verify = 0;
pool.prev_end = None;
pool.active = true;
assert!(
pool.results.iter().all(|s| s.as_ref().unwrap().started_fresh()),
"all seeded segments are Fresh guesses",
);
assert!(
pool.results[2].as_ref().unwrap().end_boundary() != &pool.fresh,
"an unterminated `\"` leaves the string context open (non-fresh end)",
);
let mut frames = 0usize;
loop {
let before = pool.next_verify;
let reruns_before = POLL_RERUNS.with(std::cell::Cell::get);
pool.poll(&mut doc);
let delta = pool.next_verify - before;
let reruns = POLL_RERUNS.with(std::cell::Cell::get) - reruns_before;
frames += 1;
assert!(delta <= POLL_VERIFY_BUDGET, "frame absorbed {delta} > verify budget");
assert!(reruns <= POLL_RERUN_BUDGET as u64, "frame ran {reruns} > rerun budget");
if pool.next_verify < pool.results.len() {
assert!(pool.active, "a partial frame stays active so the subscription re-fires");
assert!(delta >= 1, "a partial frame must make progress (no live-lock)");
}
if !pool.active {
break;
}
assert!(frames <= N + 4, "must converge, not spin");
}
assert_eq!(pool.next_verify, N, "every segment verified");
assert!(!pool.active, "active clears only once the document is fully verified");
assert!(frames > 1, "verification was paced across frames, not drained in one");
}
}