use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc::{Receiver, Sender, TryRecvError};
use std::sync::Arc;
use crate::sqlite::{PageHint, RowsView, Sort, SqliteStore};
const CANCEL_CHECK_OPS: i32 = 2_000;
#[derive(Debug, Clone)]
pub struct PageReq {
pub table: String,
pub limit: i64,
pub offset: i64,
pub sort: Option<Sort>,
pub filter: String,
pub hint: Option<PageHint>,
pub known_total: Option<i64>,
}
#[derive(Debug, Clone)]
pub struct CountReq {
pub table: String,
pub filter: String,
}
#[derive(Debug, Clone)]
pub struct SearchReq {
pub table: String,
pub columns: Vec<String>,
pub term: String,
pub sort: Option<Sort>,
pub filter: String,
pub from: Option<i64>,
pub forward: bool,
}
pub struct SearchDone {
pub generation: u64,
pub result: Result<Option<(i64, i64)>, String>,
}
pub struct PageDone {
pub generation: u64,
pub result: Result<RowsView, String>,
}
pub struct CountDone {
pub generation: u64,
pub table: String,
pub filter: String,
pub result: Result<i64, String>,
}
struct Worker<Req, Out> {
jobs: Sender<(u64, Req)>,
out: Receiver<Out>,
live: Arc<AtomicU64>,
inflight: bool,
}
impl<Req, Out> Worker<Req, Out> {
fn send(&mut self, generation: u64, req: Req) {
self.live.store(generation, Ordering::Relaxed);
if self.jobs.send((generation, req)).is_ok() {
self.inflight = true;
}
}
fn poll(&mut self) -> Option<Out> {
match self.out.try_recv() {
Ok(v) => {
self.inflight = false;
Some(v)
}
Err(TryRecvError::Empty) => None,
Err(TryRecvError::Disconnected) => {
self.inflight = false;
None
}
}
}
}
pub struct Engine {
pages: Worker<PageReq, PageDone>,
counts: Worker<CountReq, CountDone>,
searches: Worker<SearchReq, SearchDone>,
next_generation: u64,
pub page_generation: u64,
}
impl Engine {
pub fn new(path: &std::path::Path) -> Engine {
Engine {
pages: spawn_worker(path.to_path_buf(), |store, generation, req: PageReq| {
PageDone {
generation,
result: store
.rows(&crate::sqlite::PageQuery {
table: &req.table,
limit: req.limit,
offset: req.offset,
sort: req.sort.as_ref(),
filter: &req.filter,
hint: req.hint.as_ref(),
known_total: req.known_total,
})
.map_err(|e| e.to_string()),
}
}),
counts: spawn_worker(path.to_path_buf(), |store, generation, req: CountReq| {
CountDone {
generation,
result: store
.count_exact(&req.table, &req.filter)
.map_err(|e| e.to_string()),
table: req.table,
filter: req.filter,
}
}),
searches: spawn_worker(path.to_path_buf(), |store, generation, req: SearchReq| {
SearchDone {
generation,
result: search(store, &req),
}
}),
next_generation: 1,
page_generation: 0,
}
}
pub fn request(&mut self, page: PageReq, count: Option<CountReq>) -> u64 {
let generation = self.next_generation;
self.next_generation += 1;
self.pages.send(generation, page);
if let Some(count) = count {
self.counts.send(generation, count);
}
generation
}
pub fn request_count(&mut self, count: CountReq) {
let generation = self.next_generation;
self.next_generation += 1;
self.counts.send(generation, count);
}
pub fn request_search(&mut self, req: SearchReq) {
let generation = self.next_generation;
self.next_generation += 1;
self.searches.send(generation, req);
}
pub fn poll_search(&mut self) -> Option<SearchDone> {
let live = self.searches.live.load(Ordering::Relaxed);
while let Some(done) = self.searches.poll() {
if done.generation == live {
return Some(done);
}
}
None
}
pub fn searching(&self) -> bool {
self.searches.inflight
}
pub fn page_inflight(&self) -> bool {
self.pages.inflight
}
pub fn count_inflight(&self) -> bool {
self.counts.inflight
}
pub fn poll_page(&mut self) -> Option<PageDone> {
let live = self.pages.live.load(Ordering::Relaxed);
while let Some(done) = self.pages.poll() {
if done.generation == live {
self.page_generation = done.generation;
return Some(done);
}
}
None
}
pub fn wait_page(&mut self, grace: std::time::Duration) -> Option<PageDone> {
let deadline = std::time::Instant::now() + grace;
loop {
if let Some(done) = self.poll_page() {
return Some(done);
}
let left = deadline.saturating_duration_since(std::time::Instant::now());
if left.is_zero() {
return None;
}
match self.pages.out.recv_timeout(left) {
Ok(done) => {
self.pages.inflight = false;
if done.generation == self.pages.live.load(Ordering::Relaxed) {
self.page_generation = done.generation;
return Some(done);
}
}
Err(_) => return None,
}
}
}
pub fn poll_count(&mut self) -> Option<CountDone> {
let live = self.counts.live.load(Ordering::Relaxed);
while let Some(done) = self.counts.poll() {
if done.generation == live {
return Some(done);
}
}
None
}
}
fn search(store: &SqliteStore, req: &SearchReq) -> Result<Option<(i64, i64)>, String> {
let query = crate::sqlite::RowQuery {
table: &req.table,
columns: &req.columns,
term: &req.term,
sort: req.sort.as_ref(),
filter: &req.filter,
};
let first = match req.from {
Some(from) => store.find_row(&query, from, req.forward),
None => store.find_row_edge(&query, req.forward),
};
let found = match first {
Err(e) => return Err(e.to_string()),
Ok(Some(r)) => Some(r),
Ok(None) => store
.find_row_edge(&query, req.forward)
.map_err(|e| e.to_string())?,
};
let Some(rowid) = found else {
return Ok(None);
};
let ordinal = store
.rowid_ordinal(&req.table, rowid, req.sort.as_ref(), &req.filter)
.unwrap_or(1);
Ok(Some((rowid, ordinal)))
}
fn spawn_worker<Req, Out, F>(path: PathBuf, work: F) -> Worker<Req, Out>
where
Req: Send + 'static,
Out: Send + 'static,
F: Fn(&SqliteStore, u64, Req) -> Out + Send + 'static,
{
let (jobs_tx, jobs_rx) = std::sync::mpsc::channel::<(u64, Req)>();
let (out_tx, out_rx) = std::sync::mpsc::channel::<Out>();
let live = Arc::new(AtomicU64::new(0));
let live_thread = Arc::clone(&live);
std::thread::spawn(move || {
let mut store = match SqliteStore::open_readonly(&path) {
Ok(s) => s,
Err(_) => return,
};
let mine = Arc::new(AtomicU64::new(0));
let (live_cb, mine_cb) = (Arc::clone(&live_thread), Arc::clone(&mine));
store.set_cancel(
CANCEL_CHECK_OPS,
Arc::new(move || live_cb.load(Ordering::Relaxed) != mine_cb.load(Ordering::Relaxed)),
);
while let Ok(job) = jobs_rx.recv() {
let (generation, req) = {
let mut latest = job;
while let Ok(next) = jobs_rx.try_recv() {
latest = next;
}
latest
};
if generation != live_thread.load(Ordering::Relaxed) {
continue;
}
mine.store(generation, Ordering::Relaxed);
if out_tx.send(work(&store, generation, req)).is_err() {
return;
}
}
});
Worker {
jobs: jobs_tx,
out: out_rx,
live,
inflight: false,
}
}