use super::client::Seat;
use super::link::LinkEnd;
use crate::state::{SnapshotCell, latest_snapshot};
use crate::watch::Repaint;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread::JoinHandle;
use std::time::Duration;
pub const ASK_PERIOD: Duration = Duration::from_millis(500);
pub struct Asker {
seat: Seat,
end: LinkEnd,
snap: SnapshotCell,
state_root: PathBuf,
repaint: Arc<dyn Repaint>,
}
impl Asker {
pub fn new(
seat: Seat,
end: LinkEnd,
snap: SnapshotCell,
state_root: PathBuf,
repaint: Arc<dyn Repaint>,
) -> Self {
Self {
seat,
end,
snap,
state_root,
repaint,
}
}
pub fn pass(&mut self) -> usize {
self.seat_window();
let mut landed = 0;
for question in self.end.standing() {
let answer = self.seat.answered(&question);
if !self.end.publish(&question, answer) {
break;
}
landed += 1;
}
if landed > 0 {
self.repaint.request();
}
landed
}
fn seat_window(&self) {
let client = crate::registry::window();
let snap = latest_snapshot(&self.snap);
let seated = crate::registry::registered(&self.state_root, &client);
for workspace in &snap.workspaces {
let name = snap.ws_name(&workspace.path);
if !seated.contains(&name) {
let _ = crate::registry::register(&self.state_root, &client, &name);
}
}
}
pub fn start(mut self) -> AskerThread {
let stop = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&stop);
let handle = std::thread::spawn(move || {
while !flag.load(Ordering::Relaxed) {
self.pass();
std::thread::park_timeout(ASK_PERIOD);
}
});
AskerThread {
stop,
handle: Some(handle),
}
}
}
pub struct AskerThread {
stop: Arc<AtomicBool>,
handle: Option<JoinHandle<()>>,
}
impl Drop for AskerThread {
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;