use std::path::PathBuf;
use std::sync::mpsc::{Receiver, Sender, channel};
use std::thread;
use std::time::{Duration, Instant};
#[derive(Clone, Copy, PartialEq, Debug)]
pub(crate) enum Kind {
Names,
Content,
}
struct Job {
seq: u64,
kind: Kind,
root: PathBuf,
query: String,
hidden: bool,
}
pub(crate) struct Done {
seq: u64,
pub(crate) hits: Vec<PathBuf>,
pub(crate) err: Option<String>,
}
pub(crate) struct Search {
tx: Sender<Job>,
rx: Receiver<Done>,
seq: u64,
inflight: bool,
}
impl Search {
pub(crate) fn spawn() -> Self {
let (tx, job_rx) = channel::<Job>();
let (done_tx, rx) = channel::<Done>();
thread::spawn(move || worker(job_rx, done_tx));
Search {
tx,
rx,
seq: 0,
inflight: false,
}
}
pub(crate) fn request(&mut self, kind: Kind, root: PathBuf, query: String, hidden: bool) {
self.seq += 1;
self.inflight = self
.tx
.send(Job {
seq: self.seq,
kind,
root,
query,
hidden,
})
.is_ok();
}
pub(crate) fn cancel(&mut self) {
self.seq += 1;
self.inflight = false;
}
pub(crate) fn take_fresh(&mut self) -> Option<Done> {
let mut fresh = None;
while let Ok(d) = self.rx.try_recv() {
if d.seq == self.seq {
fresh = Some(d);
}
}
if fresh.is_some() {
self.inflight = false;
}
fresh
}
pub(crate) fn pending(&self) -> bool {
self.inflight
}
pub(crate) fn wait(&mut self, cap: Duration) -> Option<Done> {
let deadline = Instant::now() + cap;
while self.inflight {
let left = deadline.checked_duration_since(Instant::now())?;
match self.rx.recv_timeout(left) {
Ok(d) if d.seq == self.seq => {
self.inflight = false;
return Some(d);
}
Ok(_) => {} Err(_) => return None, }
}
None
}
}
fn worker(jobs: Receiver<Job>, out: Sender<Done>) {
while let Ok(mut job) = jobs.recv() {
while let Ok(newer) = jobs.try_recv() {
job = newer;
}
if out.send(run(job)).is_err() {
return; }
}
}
fn run(job: Job) -> Done {
match job.kind {
Kind::Names => Done {
seq: job.seq,
hits: cdt_search::find_names(&job.root, &job.query, job.hidden),
err: None,
},
Kind::Content => match cdt_search::grep(&job.root, &job.query, job.hidden) {
Ok(hits) => Done {
seq: job.seq,
hits,
err: None,
},
Err(e) => Done {
seq: job.seq,
hits: Vec::new(),
err: Some(format!("rg unavailable: {e}")),
},
},
}
}
impl std::fmt::Debug for Done {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Done")
.field("seq", &self.seq)
.field("hits", &self.hits.len())
.field("err", &self.err)
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn root() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR"))
}
#[test]
fn a_requested_search_comes_back_with_hits() {
let mut s = Search::spawn();
assert!(!s.pending());
s.request(Kind::Names, root(), "lib.rs".into(), false);
assert!(s.pending());
let done = s.wait(Duration::from_secs(10)).expect("a result");
assert!(done.hits.iter().any(|p| p.ends_with("lib.rs")), "{done:?}",);
assert!(!s.pending());
}
#[test]
fn a_result_for_a_superseded_query_is_discarded() {
let mut s = Search::spawn();
s.request(Kind::Names, root(), "lib.rs".into(), false);
s.cancel();
thread::sleep(std::time::Duration::from_millis(300));
assert!(s.take_fresh().is_none(), "stale result leaked through");
assert!(!s.pending());
}
#[test]
fn take_fresh_keeps_only_the_newest_answer() {
let mut s = Search::spawn();
for q in ["l", "li", "lib"] {
s.request(Kind::Names, root(), q.into(), false);
}
let last = s.seq;
let done = s.wait(Duration::from_secs(10)).expect("a result");
assert_eq!(done.seq, last, "answered a stale query");
}
#[test]
fn a_missing_rg_is_reported_rather_than_silently_empty() {
let mut s = Search::spawn();
s.request(Kind::Content, root(), "Search".into(), false);
let done = s.wait(Duration::from_secs(10)).expect("a result");
match done.err {
Some(e) => assert!(e.contains("rg"), "{e}"),
None => assert!(!done.hits.is_empty(), "rg found nothing in its own source"),
}
}
}