use crate::{
RoughPlannerWindow,
gui::rough_plan::{host::on_host, saved::panic_message},
};
use indicatrix_vault::db::sqlite::Database;
use slint::Weak;
use std::{
cell::Cell,
collections::BTreeMap,
panic::{AssertUnwindSafe, catch_unwind},
sync::{Arc, Mutex, PoisonError, mpsc},
};
use tracing::warn;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) enum Change {
One {
id: i64,
before: String,
excluded: bool,
},
All,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) enum Job {
Read,
Write {
ids: Vec<i64>,
excluded: bool,
change: Change,
},
}
impl Job {
#[must_use]
pub(super) const fn is_write(&self) -> bool {
matches!(self, Self::Write { .. })
}
}
pub(super) type Listed = Result<BTreeMap<i64, String>, String>;
#[derive(Debug, PartialEq, Eq)]
pub(super) struct Answer {
pub(super) write: Option<Result<Change, String>>,
pub(super) listed: Listed,
}
fn write_failure(detail: &str) -> String {
format!("Could not change the planner exclusion: {detail}")
}
fn read_excluded(db: &Mutex<Database>) -> Listed {
let db = db.lock().unwrap_or_else(PoisonError::into_inner);
let ids: Vec<i64> = db
.planner_excluded_ids()
.map_err(|e| format!("{e:#}"))?
.into_iter()
.collect();
db.entry_titles_for(&ids).map_err(|e| format!("{e:#}"))
}
fn write_marks(db: &Mutex<Database>, ids: &[i64], excluded: bool) -> Result<(), String> {
let db = db.lock().unwrap_or_else(PoisonError::into_inner);
ids.iter().try_for_each(|&id| {
db.set_planner_excluded(id, excluded)
.map_err(|e| write_failure(&format!("{e:#}")))
})
}
pub(super) fn run_job(db: &Mutex<Database>, job: Job) -> Answer {
let write = match job {
Job::Read => None,
Job::Write {
ids,
excluded,
change,
} => Some(write_marks(db, &ids, excluded).map(|()| change)),
};
Answer {
write,
listed: read_excluded(db),
}
}
fn panicked(was_write: bool, detail: &str) -> Answer {
warn!("Rough planner: an exclusion job panicked: {detail}");
let message = format!("the exclusions thread stopped unexpectedly ({detail})");
Answer {
write: was_write.then(|| Err(write_failure(&message))),
listed: Err(message),
}
}
fn spawn_queue(
db: Arc<Mutex<Database>>,
deliver: impl Fn(Answer) + Send + 'static,
) -> Option<mpsc::Sender<Job>> {
let (jobs, inbox) = mpsc::channel::<Job>();
let spawned = std::thread::Builder::new()
.name("rough-exclusions".to_string())
.spawn(move || {
for job in inbox {
let was_write = job.is_write();
let answer = catch_unwind(AssertUnwindSafe(|| run_job(&db, job)))
.unwrap_or_else(|payload| panicked(was_write, &panic_message(&*payload)));
deliver(answer);
}
});
match spawned {
Ok(_handle) => Some(jobs),
Err(error) => {
warn!("Rough planner: could not start the exclusions thread: {error}");
None
}
}
}
pub(in crate::gui::rough_plan) struct ExclusionWorker {
jobs: Option<mpsc::Sender<Job>>,
pending_writes: Cell<usize>,
}
impl ExclusionWorker {
#[must_use]
pub(in crate::gui::rough_plan) fn new(
window: Weak<RoughPlannerWindow>,
db: Arc<Mutex<Database>>,
) -> Self {
let jobs = spawn_queue(db, move |answer| {
let _ = window.upgrade_in_event_loop(move |_window| {
on_host(|host| super::apply_answer(host, answer));
});
});
Self {
jobs,
pending_writes: Cell::new(0),
}
}
pub(super) fn submit(&self, job: Job) -> bool {
let is_write = job.is_write();
let sent = self
.jobs
.as_ref()
.is_some_and(|jobs| jobs.send(job).is_ok());
if sent && is_write {
self.pending_writes.set(self.pending_writes.get() + 1);
}
sent
}
#[must_use]
pub(super) const fn writes_pending(&self) -> bool {
self.pending_writes.get() > 0
}
pub(super) fn write_answered(&self) {
self.pending_writes
.set(self.pending_writes.get().saturating_sub(1));
}
}
#[cfg(test)]
mod tests {
use super::*;
use indicatrix_vault::model::entry::FacetingDiagramEntry;
use std::time::Duration;
const WAIT: Duration = Duration::from_secs(10);
fn library(titles: &[&str]) -> (Arc<Mutex<Database>>, Vec<i64>) {
let db = Database::new(Some(":memory:")).expect("an in-memory library opens");
let ids = titles
.iter()
.enumerate()
.map(|(index, title)| {
db.save_diagram_entry(
&FacetingDiagramEntry {
title: (*title).to_string(),
url: format!("local://{title}-{index}.asc"),
design_id: String::new(),
},
"local-import",
)
.expect("the design is saved")
})
.collect();
(Arc::new(Mutex::new(db)), ids)
}
fn one(id: i64, excluded: bool) -> Change {
Change::One {
id,
before: format!("Design #{id}"),
excluded,
}
}
fn write(ids: &[i64], excluded: bool, change: Change) -> Job {
Job::Write {
ids: ids.to_vec(),
excluded,
change,
}
}
#[test]
fn a_read_lists_the_excluded_designs_with_their_titles() {
let (db, ids) = library(&["Oval", "Round"]);
assert_eq!(
run_job(&db, Job::Read),
Answer {
write: None,
listed: Ok(BTreeMap::new()),
}
);
db.lock()
.expect("lock")
.set_planner_excluded(ids[1], true)
.expect("excluded");
let answer = run_job(&db, Job::Read);
assert_eq!(answer.write, None);
assert_eq!(
answer.listed,
Ok(BTreeMap::from([(ids[1], "Round".to_string())]))
);
}
#[test]
fn a_write_marks_the_designs_and_answers_with_the_list_after_it() {
let (db, ids) = library(&["Oval", "Round", "Cushion"]);
let change = one(ids[0], true);
let answer = run_job(&db, write(&[ids[0]], true, change.clone()));
assert_eq!(answer.write, Some(Ok(change)));
assert_eq!(
answer.listed,
Ok(BTreeMap::from([(ids[0], "Oval".to_string())]))
);
run_job(&db, write(&[ids[1]], true, one(ids[1], true)));
let answer = run_job(&db, write(&[ids[0], ids[1]], false, Change::All));
assert_eq!(answer.write, Some(Ok(Change::All)));
assert_eq!(answer.listed, Ok(BTreeMap::new()));
}
#[test]
fn a_failed_write_says_so_and_the_list_is_still_read() {
let (db, ids) = library(&["Oval"]);
run_job(&db, write(&[ids[0]], true, one(ids[0], true)));
let answer = run_job(&db, write(&[ids[0] + 100], true, one(ids[0] + 100, true)));
let message = answer
.write
.expect("a write job answers")
.expect_err("the library refuses a design it does not have");
assert!(
message.starts_with("Could not change the planner exclusion:"),
"{message}"
);
assert_eq!(
answer.listed,
Ok(BTreeMap::from([(ids[0], "Oval".to_string())])),
"the mark that is there stays in the list"
);
}
#[test]
fn a_panic_is_answered_as_a_failed_write_or_a_failed_read() {
let failed_write = panicked(true, "boom");
let message = failed_write.write.expect("a write answers").unwrap_err();
assert!(message.contains("boom"), "{message}");
assert!(message.starts_with("Could not change the planner exclusion:"));
assert!(failed_write.listed.is_err());
let read = panicked(false, "boom");
assert_eq!(read.write, None);
assert!(read.listed.unwrap_err().contains("boom"));
}
#[test]
fn jobs_are_answered_in_the_order_they_were_asked_for() {
let (db, ids) = library(&["Oval", "Round"]);
let (answers_tx, answers_rx) = mpsc::channel::<Answer>();
let jobs = spawn_queue(Arc::clone(&db), move |answer| {
let _ = answers_tx.send(answer);
})
.expect("the thread starts");
jobs.send(write(&[ids[0]], true, one(ids[0], true)))
.expect("sent");
jobs.send(Job::Read).expect("sent");
jobs.send(write(&[ids[0]], false, one(ids[0], false)))
.expect("sent");
jobs.send(Job::Read).expect("sent");
let listed = |answer: Answer| answer.listed.expect("readable").len();
let sizes: Vec<usize> = (0..4)
.map(|_| listed(answers_rx.recv_timeout(WAIT).expect("an answer")))
.collect();
assert_eq!(sizes, vec![1, 1, 0, 0]);
drop(jobs);
assert!(answers_rx.recv_timeout(WAIT).is_err());
}
#[test]
fn a_job_that_cannot_be_queued_is_reported() {
let (jobs, inbox) = mpsc::channel::<Job>();
drop(inbox);
let worker = ExclusionWorker {
jobs: Some(jobs),
pending_writes: Cell::new(0),
};
assert!(!worker.submit(Job::Read));
assert!(!worker.submit(write(&[1], true, Change::All)));
assert!(
!worker.writes_pending(),
"a write that was not queued is not awaited"
);
let without_thread = ExclusionWorker {
jobs: None,
pending_writes: Cell::new(0),
};
assert!(!without_thread.submit(Job::Read));
}
#[test]
fn a_write_is_pending_until_its_answer_arrives() {
let (jobs, _inbox) = mpsc::channel::<Job>();
let worker = ExclusionWorker {
jobs: Some(jobs),
pending_writes: Cell::new(0),
};
assert!(worker.submit(Job::Read));
assert!(!worker.writes_pending(), "a read is not a write");
assert!(worker.submit(write(&[1], true, Change::All)));
assert!(worker.writes_pending());
worker.write_answered();
assert!(!worker.writes_pending());
worker.write_answered();
assert!(
!worker.writes_pending(),
"an extra answer does not go below none"
);
}
}