use super::{model::SettingsFile, store};
use std::{
cell::RefCell,
path::PathBuf,
sync::{Arc, Mutex, Weak, mpsc},
thread,
time::{Duration, Instant},
};
use tracing::warn;
const DEBOUNCE: Duration = Duration::from_millis(600);
const POLL_INTERVAL: Duration = Duration::from_millis(100);
const FLUSH_WAIT: Duration = Duration::from_secs(3);
enum Msg {
Changed(SettingsFile),
FlushNow(SettingsFile, mpsc::Sender<()>),
Shutdown,
}
thread_local! {
static INSTALLED: RefCell<Weak<SettingsPersister>> = const { RefCell::new(Weak::new()) };
}
pub type DiskOverride = Arc<dyn Fn(&mut SettingsFile) + Send + Sync>;
#[derive(Clone)]
pub struct SettingsPersister {
sender: mpsc::Sender<Msg>,
current: Arc<Mutex<SettingsFile>>,
disk_override: Arc<Mutex<Option<DiskOverride>>>,
}
impl SettingsPersister {
#[must_use]
pub fn spawn(path: PathBuf, initial: SettingsFile) -> Self {
let (sender, receiver) = mpsc::channel::<Msg>();
let current = Arc::new(Mutex::new(initial));
thread::spawn(move || worker_loop(&path, &receiver));
Self {
sender,
current,
disk_override: Arc::new(Mutex::new(None)),
}
}
pub fn update(&self, f: impl FnOnce(&mut SettingsFile)) {
let mut guard = self
.current
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
f(&mut guard);
let snapshot = guard.clone();
drop(guard);
let _ = self.sender.send(Msg::Changed(self.on_disk(snapshot)));
}
pub fn set_disk_override(&self, rewrite: Option<DiskOverride>) {
*self
.disk_override
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = rewrite;
let _ = self
.sender
.send(Msg::Changed(self.on_disk(self.snapshot())));
}
fn on_disk(&self, mut snapshot: SettingsFile) -> SettingsFile {
let rewrite = self
.disk_override
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(rewrite) = rewrite {
rewrite(&mut snapshot);
}
snapshot
}
#[must_use]
pub fn snapshot(&self) -> SettingsFile {
self.current
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub fn flush(&self) {
let snapshot = self.on_disk(self.snapshot());
let (ack_sender, ack_receiver) = mpsc::channel();
if self
.sender
.send(Msg::FlushNow(snapshot, ack_sender))
.is_ok()
{
let _ = ack_receiver.recv_timeout(FLUSH_WAIT);
}
}
pub fn install_for_this_thread(handle: &Arc<Self>) {
INSTALLED.with(|slot| *slot.borrow_mut() = Arc::downgrade(handle));
}
#[must_use]
pub fn installed_for_this_thread() -> Option<Arc<Self>> {
INSTALLED.with(|slot| slot.borrow().upgrade())
}
}
impl Drop for SettingsPersister {
fn drop(&mut self) {
let _ = self.sender.send(Msg::Shutdown);
}
}
fn worker_loop(path: &std::path::Path, receiver: &mpsc::Receiver<Msg>) {
let mut pending: Option<SettingsFile> = None;
let mut last_change = Instant::now();
loop {
match receiver.recv_timeout(POLL_INTERVAL) {
Ok(Msg::Changed(settings)) => {
pending = Some(settings);
last_change = Instant::now();
}
Ok(Msg::FlushNow(settings, ack)) => {
write_settings(path, &settings);
pending = None;
let _ = ack.send(());
}
Ok(Msg::Shutdown) => {
if let Some(settings) = pending.take() {
write_settings(path, &settings);
}
break;
}
Err(mpsc::RecvTimeoutError::Timeout) => {
if let Some(settings) = &pending
&& last_change.elapsed() >= DEBOUNCE
{
write_settings(path, settings);
pending = None;
}
}
Err(mpsc::RecvTimeoutError::Disconnected) => break,
}
}
}
fn write_settings(path: &std::path::Path, settings: &SettingsFile) {
if let Err(e) = store::save(path, settings) {
warn!("Failed to save settings to {}: {e}", path.display());
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU32, Ordering};
static COUNTER: AtomicU32 = AtomicU32::new(0);
fn temp_settings_path(tag: &str) -> PathBuf {
let n = COUNTER.fetch_add(1, Ordering::Relaxed);
let dir = std::env::temp_dir().join(format!(
"indicatrix-cut-persist-test-{tag}-{}-{n}",
std::process::id()
));
std::fs::create_dir_all(&dir).unwrap();
dir.join("settings.toml")
}
#[test]
fn flush_writes_immediately_without_waiting_for_the_debounce() {
let path = temp_settings_path("flush");
let persister = SettingsPersister::spawn(path.clone(), SettingsFile::default());
persister.update(|s| s.settings.exposure = 2.5);
persister.flush();
let deadline = Instant::now() + Duration::from_secs(5);
loop {
if let Ok(contents) = std::fs::read_to_string(&path)
&& contents.contains("exposure = 2.5")
{
break;
}
assert!(Instant::now() < deadline, "flush did not persist in time");
thread::sleep(Duration::from_millis(20));
}
}
#[test]
fn snapshot_reflects_update_immediately_even_before_the_write_lands() {
let path = temp_settings_path("snapshot");
let persister = SettingsPersister::spawn(path, SettingsFile::default());
persister.update(|s| s.settings.max_bounces = 42);
assert_eq!(persister.snapshot().settings.max_bounces, 42);
}
#[test]
fn a_disk_override_keeps_loaned_values_out_of_the_file_but_not_out_of_memory() {
let path = temp_settings_path("disk-override");
let persister = SettingsPersister::spawn(path.clone(), SettingsFile::default());
persister.update(|s| s.settings.exposure = 1.5);
persister.flush();
assert_eq!(store::load_or_default(&path).settings.exposure, 1.5);
persister.set_disk_override(Some(Arc::new(|file: &mut SettingsFile| {
file.settings.exposure = 1.5;
})));
persister.update(|s| {
s.settings.exposure = 2.5;
s.settings.max_bounces = 42;
});
assert_eq!(
persister.snapshot().settings.exposure,
2.5,
"the controls still read what they set"
);
persister.flush();
let on_disk = store::load_or_default(&path).settings;
assert_eq!(
on_disk.exposure, 1.5,
"the normal exposure is what the file keeps"
);
assert_eq!(
on_disk.max_bounces, 42,
"a setting the override does not name is saved"
);
persister.set_disk_override(None);
persister.update(|s| s.settings.exposure = 3.5);
persister.flush();
assert_eq!(store::load_or_default(&path).settings.exposure, 3.5);
}
#[test]
fn a_disk_override_also_applies_to_the_debounced_write() {
let path = temp_settings_path("disk-override-debounced");
let persister = SettingsPersister::spawn(path.clone(), SettingsFile::default());
persister.set_disk_override(Some(Arc::new(|file: &mut SettingsFile| {
file.settings.exposure = 1.5;
})));
persister.update(|s| s.settings.exposure = 2.5);
let deadline = Instant::now() + Duration::from_secs(10);
let written = loop {
if let Ok(contents) = std::fs::read_to_string(&path)
&& contents.contains("exposure = ")
{
break contents;
}
assert!(Instant::now() < deadline, "the debounced write never came");
thread::sleep(Duration::from_millis(50));
};
assert!(written.contains("exposure = 1.5"), "{written}");
assert!(!written.contains("exposure = 2.5"), "{written}");
}
#[test]
fn a_recent_file_recorded_through_the_persister_survives_flush() {
let path = temp_settings_path("recent-survives");
let persister = SettingsPersister::spawn(path.clone(), SettingsFile::default());
persister.update(|s| {
s.settings
.record_recent_native_file("saved-this-session.indicatrix.toml".to_owned());
});
persister.flush();
let on_disk = store::load_or_default(&path);
assert_eq!(
on_disk.settings.recent_native_files,
vec!["saved-this-session.indicatrix.toml".to_owned()]
);
}
#[test]
fn flush_overwrites_a_file_written_behind_the_persisters_back() {
let path = temp_settings_path("behind-the-back");
let persister = SettingsPersister::spawn(path.clone(), SettingsFile::default());
let mut direct = SettingsFile::default();
direct
.settings
.record_recent_native_file("written-directly.indicatrix.toml".to_owned());
store::save(&path, &direct).unwrap();
persister.flush();
assert_eq!(
store::load_or_default(&path).settings.recent_native_files,
Vec::<String>::new()
);
}
#[test]
fn the_installed_handle_is_per_thread_and_shares_the_persisters_state() {
let path = temp_settings_path("installed");
let persister = Arc::new(SettingsPersister::spawn(path, SettingsFile::default()));
SettingsPersister::install_for_this_thread(&persister);
let installed = SettingsPersister::installed_for_this_thread().expect("installed");
assert!(Arc::ptr_eq(&installed, &persister));
installed.update(|s| s.settings.max_bounces = 17);
assert_eq!(persister.snapshot().settings.max_bounces, 17);
let elsewhere = thread::spawn(|| SettingsPersister::installed_for_this_thread().is_none())
.join()
.unwrap();
assert!(elsewhere);
}
}