use crate::state::SnapshotCell;
use crate::ui_state::{Clock, UiState, content_hash};
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread::JoinHandle;
use std::time::Duration;
use super::consume::consume;
use super::deposit;
use super::dispatch::Deps;
const CONSUMER_POLL: Duration = Duration::from_millis(250);
pub struct ConsumerCtx {
pub lernie: crate::cli_outbound::Cli,
pub bl: crate::cli_outbound::Cli,
pub state_root: PathBuf,
pub home: PathBuf,
pub yog_data_root: PathBuf,
pub balls_state_root: PathBuf,
pub yog_binary: PathBuf,
pub world: crate::xdg::Env,
pub ui_path: PathBuf,
pub cell: SnapshotCell,
pub clock: Arc<dyn Clock>,
}
impl ConsumerCtx {
pub fn pass(&self) -> usize {
if deposit::pending(&self.state_root).is_empty() {
return 0;
}
let ts = self.clock.stamp();
let now_unix: i64 = ts.parse().unwrap_or(0);
let deps = Deps {
lernie: self.lernie.clone(),
bl: self.bl.clone(),
state_root: self.state_root.clone(),
yog_binary: self.yog_binary.clone(),
world: self.world.clone(),
home: self.home.clone(),
yog_data_root: self.yog_data_root.clone(),
balls_state_root: self.balls_state_root.clone(),
snapshot: crate::state::latest_snapshot(&self.cell),
mint_seed: content_hash(ts.as_bytes()),
};
let mut ui = UiState::open(self.ui_path.clone());
consume(&deps, &mut ui, &ts, now_unix)
}
}
pub struct Consumer {
stop: Arc<AtomicBool>,
handle: Option<JoinHandle<()>>,
}
impl Consumer {
pub fn spawn(ctx: ConsumerCtx) -> Self {
let stop = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&stop);
let handle = std::thread::spawn(move || {
while !flag.load(Ordering::Relaxed) {
ctx.pass();
std::thread::park_timeout(CONSUMER_POLL);
}
});
Self {
stop,
handle: Some(handle),
}
}
}
impl Drop for Consumer {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(handle) = self.handle.take() {
handle.thread().unpark();
let _ = handle.join();
}
}
}
#[cfg(test)]
mod tests;