use std::path::{Path, PathBuf};
use crate::config::Config;
use crate::db::connection::Database;
use super::*;
pub fn try_watch_lock(db_path: &Path) -> std::io::Result<Option<std::fs::File>> {
let file = watch_lock_file(db_path)?;
match file.try_lock() {
Ok(()) => Ok(Some(file)),
Err(std::fs::TryLockError::WouldBlock) => Ok(None),
Err(std::fs::TryLockError::Error(e)) => Err(e),
}
}
pub fn watch_lock(db_path: &Path) -> std::io::Result<std::fs::File> {
let file = watch_lock_file(db_path)?;
file.lock()?;
Ok(file)
}
fn watch_lock_file(db_path: &Path) -> std::io::Result<std::fs::File> {
std::fs::File::options()
.create(true)
.truncate(false)
.write(true)
.open(db_path.with_file_name("watch.lock"))
}
pub fn spawn_library_watch(
db_path: std::path::PathBuf,
on_state: impl Fn(bool) + Send + Sync + 'static,
) -> Option<std::thread::JoinHandle<()>> {
use std::collections::BTreeSet;
use std::time::{Duration, Instant};
use notify::{RecursiveMode, Watcher};
use crate::index::scanner::{self, ScanOptions};
use crate::index::watch::{WatchedRoot, scan_target};
const SETTLE: Duration = Duration::from_secs(5);
const CHECK: Duration = Duration::from_secs(30);
const MAX_DIRS: usize = 200;
const RESCAN: Duration = Duration::from_secs(15 * 60);
std::thread::Builder::new()
.name("koan-library-watch".into())
.spawn(move || {
let _lock = match try_watch_lock(&db_path) {
Ok(Some(lock)) => Some(lock),
Ok(None) => {
log::info!(
"another koan is watching this library; serving, and waiting to take over"
);
let lock = watch_lock(&db_path);
if lock.is_ok() {
log::info!("library watch: taken over");
}
lock.inspect_err(|e| log::warn!("library watch: lock: {e}"))
.ok()
}
Err(e) => {
log::warn!("library watch: no lock beside the database: {e}");
None
}
};
let scan = |reason: &str, folders: &[PathBuf], dirs: Option<&[PathBuf]>| {
if folders.is_empty() {
return;
}
let Ok(db) = Database::open_existing(&db_path) else {
return;
};
if crate::db::pool::understood(&db.conn).is_err() {
return;
}
on_state(true);
let result = match dirs {
Some(dirs) => {
scanner::scan_dirs(&db, folders, dirs, ScanOptions::default(), None)
}
None => scanner::full_scan(&db, folders, ScanOptions::default(), None),
};
on_state(false);
log::info!(
"{reason} scan: {} added, {} updated, {} removed, {} unchanged",
result.added,
result.updated,
result.removed,
result.skipped
);
};
let folders = || Config::cached().library.folders.clone();
let (tx, rx) = std::sync::mpsc::channel();
let Ok(mut watcher) = notify::recommended_watcher(move |event| {
let _ = tx.send(event);
}) else {
log::warn!("could not watch the library folders");
return;
};
let mut roots: Vec<WatchedRoot> = Vec::new();
let mut rewatch = |roots: &mut Vec<WatchedRoot>| {
let wanted: Vec<WatchedRoot> = folders()
.iter()
.filter_map(|f| WatchedRoot::resolve(f))
.collect();
roots.retain(|root| {
let keep = wanted.contains(root);
if !keep {
let _ = watcher.unwatch(&root.path);
}
keep
});
let mut fresh = Vec::new();
for root in wanted {
if roots.contains(&root) {
continue;
}
match watcher.watch(&root.path, RecursiveMode::Recursive) {
Ok(()) => {
fresh.push(root.path.clone());
roots.push(root);
}
Err(e) => log::warn!("could not watch {}: {e}", root.path.display()),
}
}
fresh
};
std::thread::sleep(Duration::from_secs(3));
rewatch(&mut roots);
scan("startup", &folders(), None);
let mut dirs = BTreeSet::new();
let mut everything = false;
let mut settle_at: Option<Instant> = None;
let mut check_at = Instant::now() + CHECK;
let mut rescan_at = Instant::now() + RESCAN;
loop {
crate::quiet::wait_until_awake();
let now = Instant::now();
let wake = settle_at
.map_or(check_at, |at| at.min(check_at))
.min(rescan_at);
match rx.recv_timeout(wake.saturating_duration_since(now)) {
Ok(Ok(event)) if event.need_rescan() => {
everything = true;
settle_at = Some(Instant::now() + SETTLE);
}
Ok(Ok(event)) => {
let mut heard = false;
for dir in event
.paths
.iter()
.filter_map(|p| scan_target(&event.kind, p, &roots))
{
dirs.insert(dir);
heard = true;
}
if heard {
settle_at = Some(Instant::now() + SETTLE);
}
}
Ok(Err(e)) => log::debug!("library watch: {e}"),
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {}
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => break,
}
let now = Instant::now();
if settle_at.is_some_and(|at| now >= at) {
let changed =
scanner::minimal_dirs(std::mem::take(&mut dirs).into_iter().collect());
if everything || changed.len() > MAX_DIRS {
scan("watched change", &folders(), None);
rescan_at = Instant::now() + RESCAN;
} else {
scan("watched change", &folders(), Some(&changed));
}
everything = false;
settle_at = None;
}
if now >= rescan_at {
scan("periodic", &folders(), None);
rescan_at = Instant::now() + RESCAN;
}
if now >= check_at {
let fresh = rewatch(&mut roots);
if !fresh.is_empty() {
scan("newly watched", &fresh, None);
}
check_at = Instant::now() + CHECK;
}
}
})
.ok()
}
pub fn spawn_auto_sync(
db_path: std::path::PathBuf,
on_state: impl Fn(bool) + Send + 'static,
on_progress: impl Fn(crate::remote::sync::SyncProgress) + Send + Sync + 'static,
) -> Option<std::thread::JoinHandle<()>> {
std::thread::Builder::new()
.name("koan-auto-sync".into())
.spawn(move || {
std::thread::sleep(std::time::Duration::from_secs(5));
loop {
crate::quiet::wait_until_awake();
let cfg = Config::load().unwrap_or_default();
if !cfg.remote.enabled || !cfg.remote.auto_sync {
std::thread::sleep(std::time::Duration::from_secs(60));
continue;
}
if let Some(client) = subsonic_client(&cfg)
&& let Ok(db) = Database::open_existing(&db_path)
{
on_state(true);
match sync_remote(
&db,
&client,
Walk::IfChanged,
&cfg.remote.url,
&cfg.remote.username,
&on_progress,
) {
Ok(s) => log::info!(
"auto sync: {} artists, {} albums, {} tracks ({} albums failed); \
favourites {}↑ {}↓; playlists {}↓ {}↑",
s.library.artists_synced,
s.library.albums_synced,
s.library.tracks_synced,
s.library.albums_failed,
s.favourites.pushed,
s.favourites.imported,
s.playlists.pulled,
s.playlists.pushed,
),
Err(e) => log::warn!("auto sync failed: {e}"),
}
on_state(false);
}
let mins = cfg.remote.auto_sync_interval_mins;
if mins == 0 {
return;
}
loop {
std::thread::sleep(std::time::Duration::from_secs(mins * 60));
if !crate::remote::profile::current().is_some_and(|p| p.links()) {
break;
}
}
}
})
.ok()
}
#[cfg(test)]
mod watch_lock_tests {
use super::*;
#[test]
fn one_watcher_per_database_and_the_next_takes_over() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("koan.db");
let first = try_watch_lock(&db).unwrap().expect("the first takes it");
assert!(try_watch_lock(&db).unwrap().is_none(), "the second waits");
assert!(dir.path().join("watch.lock").exists());
drop(first);
assert!(try_watch_lock(&db).unwrap().is_some(), "and then takes it");
}
#[test]
fn a_waiting_watcher_wakes_when_the_lock_is_released() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("koan.db");
let first = try_watch_lock(&db).unwrap().unwrap();
let (tx, rx) = std::sync::mpsc::channel();
let waiting = db.clone();
std::thread::spawn(move || {
let lock = watch_lock(&waiting);
let _ = tx.send(lock.is_ok());
});
assert!(
rx.recv_timeout(std::time::Duration::from_millis(200))
.is_err(),
"it waits while the lock is held"
);
drop(first);
assert_eq!(rx.recv_timeout(std::time::Duration::from_secs(5)), Ok(true));
}
#[test]
fn databases_in_different_directories_do_not_share_a_lock() {
let (a, b) = (tempfile::tempdir().unwrap(), tempfile::tempdir().unwrap());
let _a = try_watch_lock(&a.path().join("koan.db")).unwrap().unwrap();
assert!(try_watch_lock(&b.path().join("koan.db")).unwrap().is_some());
}
}