use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use crate::config::Config;
use crate::db::connection::Database;
use crate::db::queries;
use crate::db::queries::shares::{ShareKind, Slice};
use crate::player::commands::PlayerCommand;
use crate::player::state::{ItemState, PlaylistItem, QueueItemId, SharedPlayerState};
use crate::remote::client::{Credential, SubsonicAuth, SubsonicClient, SubsonicError};
use crate::remote::download::DownloadError;
pub fn remote_credential(cfg: &Config) -> Option<Credential> {
if !cfg.remote.api_key.is_empty() {
return Some(Credential::ApiKey(cfg.remote.api_key.clone()));
}
(!cfg.remote.password.is_empty()).then(|| Credential::Password(cfg.remote.password.clone()))
}
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()
}
#[derive(Debug, Clone, Copy, Default)]
pub struct RebuildSummary {
pub tracks: u64,
pub albums: u64,
pub artists: u64,
}
pub fn rebuild_index(db: &Database) -> Result<RebuildSummary, crate::db::connection::DbError> {
let count = |sql: &str| -> u64 {
db.conn
.query_row(sql, [], |r| r.get::<_, i64>(0))
.unwrap_or(0) as u64
};
let summary = RebuildSummary {
tracks: count("SELECT COUNT(*) FROM tracks"),
albums: count("SELECT COUNT(*) FROM albums"),
artists: count("SELECT COUNT(*) FROM artists"),
};
db.conn.execute_batch(
"BEGIN;
DELETE FROM local_files;
DELETE FROM remote_entries;
DELETE FROM scan_cache;
UPDATE remote_servers SET library_version = NULL;
COMMIT;",
)?;
Ok(summary)
}
pub fn evict_cache(
db: &Database,
cfg: &Config,
keep: &std::collections::HashSet<i64>,
verbose: bool,
) -> u64 {
let Some(limit) = cfg.cache_limit_bytes().map(|l| l as i64) else {
return 0;
};
let mut current = match queries::total_cache_size(&db.conn) {
Ok(s) => s,
Err(e) => {
log::warn!("cache eviction: failed to query cache size: {e}");
return 0;
}
};
if current <= limit {
if verbose {
log::info!("cache within limit: {current} / {limit} bytes");
}
return 0;
}
let files = match queries::cached_files_lru(&db.conn) {
Ok(f) => f,
Err(e) => {
log::warn!("cache eviction: failed to query cached files: {e}");
return 0;
}
};
let mut gone = Vec::new();
let mut freed: i64 = 0;
for file in files.iter().filter(|f| !keep.contains(&f.track_id)) {
if current <= limit {
break;
}
match std::fs::remove_file(&file.path) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => {
log::warn!("cache eviction: failed to delete {}: {e}", file.path);
continue;
}
}
log::info!(
"evicted: {} ({} bytes{})",
file.path,
file.size,
if file.pinned { ", pinned" } else { "" }
);
gone.push(file.track_id);
current -= file.size;
freed += file.size;
}
if let Err(e) = queries::clear_cached_paths_for(&db.conn, &gone) {
log::warn!("cache eviction: failed to clear DB: {e}");
}
remove_empty_dirs(&cfg.cache_dir());
if freed > 0 {
cache_shrank(freed as u64);
log::info!("cache eviction freed {freed} bytes");
}
freed as u64
}
pub fn playback_window(
db: &Database,
limit: u64,
upcoming: &[i64],
) -> Result<usize, crate::db::connection::DbError> {
let total = queries::total_cache_size(&db.conn)?;
let evictable: std::collections::HashMap<i64, i64> = queries::cached_files_lru(&db.conn)?
.into_iter()
.filter(|f| !f.pinned)
.map(|f| (f.track_id, f.size))
.collect();
let estimates = queries::download_estimates(&db.conn, upcoming)?;
let mut used = (total - evictable.values().sum::<i64>()).max(0) as u64;
let mut counted = std::collections::HashSet::new();
for (n, id) in upcoming.iter().enumerate() {
if counted.insert(*id) {
let cost = evictable.get(id).or_else(|| estimates.get(id));
used += cost.copied().unwrap_or(0).max(0) as u64;
}
if n >= 2 && used > limit {
return Ok(n);
}
}
Ok(upcoming.len())
}
fn remove_empty_dirs(dir: &Path) {
if !dir.is_dir() {
return;
}
for entry in walkdir::WalkDir::new(dir)
.contents_first(true)
.into_iter()
.filter_map(Result::ok)
.filter(|e| e.file_type().is_dir() && e.path() != dir)
{
let _ = std::fs::remove_dir(entry.path());
}
}
static CACHE_BYTES: AtomicU64 = AtomicU64::new(0);
pub fn cache_bytes() -> u64 {
CACHE_BYTES.load(Ordering::Relaxed)
}
pub fn measure_cache(cfg: &Config) -> u64 {
let bytes = cache_size_bytes(cfg);
CACHE_BYTES.store(bytes, Ordering::Relaxed);
crate::signal::engine_changed().bump();
bytes
}
fn cache_grew(bytes: u64) {
CACHE_BYTES.fetch_add(bytes, Ordering::Relaxed);
crate::signal::engine_changed().bump();
}
pub(crate) fn cache_shrank(bytes: u64) {
let mut now = CACHE_BYTES.load(Ordering::Relaxed);
while let Err(moved) = CACHE_BYTES.compare_exchange_weak(
now,
now.saturating_sub(bytes),
Ordering::Relaxed,
Ordering::Relaxed,
) {
now = moved;
}
crate::signal::engine_changed().bump();
}
pub fn cache_size_bytes(cfg: &Config) -> u64 {
walkdir::WalkDir::new(cfg.cache_dir())
.into_iter()
.filter_map(Result::ok)
.filter(|e| e.file_type().is_file())
.filter(|e| e.path().extension().is_none_or(|ext| ext != "part"))
.filter_map(|e| e.metadata().ok())
.map(|m| m.len())
.sum()
}
pub fn tracks_under(db: &Database, folder: &Path) -> u64 {
let (lower, upper) = queries::folder_prefix_range(folder);
db.conn
.query_row(
"SELECT COUNT(*) FROM tracks WHERE path >= ?1 AND path < ?2",
[&lower, &upper],
|r| r.get::<_, i64>(0),
)
.unwrap_or(0) as u64
}
pub fn tracks_from_server(db: &Database) -> u64 {
db.conn
.query_row(
"SELECT COUNT(*) FROM tracks WHERE remote_id IS NOT NULL",
[],
|r| r.get::<_, i64>(0),
)
.unwrap_or(0) as u64
}
pub fn forget_folder(db: &Database, folder: &Path) -> Result<u64, crate::db::connection::DbError> {
let folder = &crate::index::spelling::on_disk(folder);
let (lower, upper) = queries::folder_prefix_range(folder);
crate::index::lane::cancel_under(folder);
let _lane = crate::index::lane::wait();
let tx = crate::db::queries::write_transaction(&db.conn)?;
let paths: Vec<String> = {
let mut stmt = tx.prepare("SELECT path FROM local_files WHERE path >= ?1 AND path < ?2")?;
let rows = stmt.query_map([&lower, &upper], |r| r.get(0))?;
rows.collect::<rusqlite::Result<_>>()?
};
for path in &paths {
queries::sources::remove(&tx, queries::sources::Kind::Local, path)?;
}
tx.execute(
"DELETE FROM scan_cache WHERE path >= ?1 AND path < ?2",
[&lower, &upper],
)?;
tx.commit()?;
Ok(paths.len() as u64)
}
pub fn forget_remote(db: &Database) -> Result<u64, crate::db::connection::DbError> {
let tx = crate::db::queries::write_transaction(&db.conn)?;
let ids: Vec<String> = {
let mut stmt = tx.prepare("SELECT remote_id FROM remote_entries")?;
let rows = stmt.query_map([], |r| r.get(0))?;
rows.collect::<rusqlite::Result<_>>()?
};
let mut removed = 0;
for id in &ids {
let track: i64 = tx.query_row(
"SELECT track_id FROM remote_entries WHERE remote_id = ?1",
[id],
|r| r.get(0),
)?;
queries::sources::remove(&tx, queries::sources::Kind::Remote, id)?;
let kept: bool = tx.query_row(
"SELECT EXISTS (SELECT 1 FROM tracks WHERE id = ?1)",
[track],
|r| r.get(0),
)?;
removed += u64::from(!kept);
}
tx.execute_batch(
"DELETE FROM history_outbox;
DELETE FROM favourite_outbox;
UPDATE remote_servers SET history_cursor = NULL;
DELETE FROM dsp_synced;
DELETE FROM dsp_sync_cursor;",
)?;
tx.commit()?;
Ok(removed)
}
#[derive(Debug, Clone, Copy, Default)]
pub struct CacheCleared {
pub files: u64,
pub bytes: u64,
}
pub fn clear_download_cache(db: &Database, cfg: &Config) -> CacheCleared {
let dir = cfg.cache_dir();
let mut cleared = CacheCleared::default();
for entry in walkdir::WalkDir::new(&dir)
.into_iter()
.filter_map(Result::ok)
.filter(|e| e.file_type().is_file())
{
if let Ok(meta) = entry.metadata() {
cleared.bytes += meta.len();
cleared.files += 1;
}
}
let _ = std::fs::remove_dir_all(&dir);
let _ = std::fs::create_dir_all(&dir);
let _ = queries::clear_cached_paths(&db.conn);
measure_cache(cfg);
cleared
}
pub fn clear_downloads_for(db: &Database, track_ids: &[i64]) -> CacheCleared {
let mut cleared = CacheCleared::default();
let paths = match queries::cached_paths_for(&db.conn, track_ids) {
Ok(paths) => paths,
Err(e) => {
log::warn!("could not read cached paths: {e}");
return cleared;
}
};
for path in &paths {
let size = std::fs::metadata(path).map(|m| m.len()).unwrap_or(0);
match std::fs::remove_file(path) {
Ok(()) => {
cleared.files += 1;
cleared.bytes += size;
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => log::warn!("could not remove {path}: {e}"),
}
}
if let Err(e) = queries::clear_cached_paths_for(&db.conn, track_ids) {
log::warn!("removed downloads but failed to forget them ({e})");
}
cache_shrank(cleared.bytes);
cleared
}
pub fn sweep_partial_downloads(cfg: &Config) -> CacheCleared {
let mut swept = CacheCleared::default();
for entry in walkdir::WalkDir::new(cfg.cache_dir())
.into_iter()
.filter_map(Result::ok)
.filter(|e| e.file_type().is_file())
.filter(|e| e.path().extension().is_some_and(|ext| ext == "part"))
{
let size = entry.metadata().map(|m| m.len()).unwrap_or(0);
match std::fs::remove_file(entry.path()) {
Ok(()) => {
swept.files += 1;
swept.bytes += size;
}
Err(e) => log::warn!("could not remove {}: {e}", entry.path().display()),
}
}
if swept.files > 0 {
log::info!(
"swept {} unfinished download(s), {} bytes",
swept.files,
swept.bytes
);
}
swept
}
pub fn relocate_cached_paths(db: &Database, cache_dir: &Path) -> rusqlite::Result<usize> {
let prefix = format!("{}/", cache_dir.to_string_lossy().trim_end_matches('/'));
let stale: Vec<(i64, String)> = db
.conn
.prepare(
"SELECT id, cached_path FROM tracks
WHERE cached_path IS NOT NULL AND substr(cached_path, 1, ?2) != ?1",
)?
.query_map(
rusqlite::params![prefix, prefix.chars().count() as i64],
|r| Ok((r.get(0)?, r.get(1)?)),
)?
.collect::<rusqlite::Result<_>>()?;
if stale.is_empty() {
return Ok(0);
}
let tx = crate::db::queries::write_transaction(&db.conn)?;
let mut moved = 0;
for (id, old) in &stale {
let tail: Vec<_> = Path::new(old).components().rev().take(3).collect();
if tail.len() < 3 {
continue;
}
let new = tail
.iter()
.rev()
.fold(cache_dir.to_path_buf(), |p, c| p.join(c));
if new.is_file() {
tx.execute(
"UPDATE tracks SET cached_path = ?1 WHERE id = ?2",
rusqlite::params![new.to_string_lossy(), id],
)?;
moved += 1;
}
}
tx.commit()?;
if moved > 0 {
log::info!(
"re-rooted {moved} cached path(s) under {}",
cache_dir.display()
);
}
Ok(moved)
}
pub fn requeue_cleared_downloads(state: &SharedPlayerState) {
let stale = state.reset_items_with_missing_files();
if !stale.is_empty() {
log::info!(
"{} queued tracks lost their copy — fetching again",
stale.len()
);
}
}
pub fn sync_favourite_to_remote(db: &Database, track_id: i64, star: bool) {
let cfg = Config::load().unwrap_or_default();
if !cfg.remote.enabled {
return;
}
let Ok(Some(remote_id)) = queries::track_remote_id(&db.conn, track_id) else {
log::warn!("not syncing favourite: track {track_id} has no remote id");
return;
};
push_favourite(db, &cfg, FavouriteKind::Track, remote_id, star);
}
fn push_favourite(db: &Database, cfg: &Config, kind: FavouriteKind, remote_id: String, star: bool) {
if let Err(e) = queries::queue_favourite_change(&db.conn, kind.as_str(), &remote_id, star) {
log::warn!("could not record a favourite change: {e}");
}
let Some(client) = subsonic_client(cfg) else {
log::warn!("not syncing favourite: no usable server credentials");
return;
};
std::thread::Builder::new()
.name("koan-fav-sync".into())
.spawn(
move || match send_favourite(&client, kind, &remote_id, star) {
Ok(()) => log::info!("synced favourite to remote: {kind:?} {remote_id} = {star}"),
Err(e) => log::warn!("failed to sync favourite to remote: {e}"),
},
)
.ok();
}
fn send_favourite(
client: &SubsonicClient,
kind: FavouriteKind,
remote_id: &str,
star: bool,
) -> Result<(), crate::remote::client::SubsonicError> {
match (kind, star) {
(FavouriteKind::Track, true) => client.star(remote_id),
(FavouriteKind::Track, false) => client.unstar(remote_id),
(FavouriteKind::Album, true) => client.star_album(remote_id),
(FavouriteKind::Album, false) => client.unstar_album(remote_id),
(FavouriteKind::Artist, true) => client.star_artist(remote_id),
(FavouriteKind::Artist, false) => client.unstar_artist(remote_id),
}
}
#[derive(Debug, Default)]
pub struct Synced {
pub library: crate::remote::sync::SyncResult,
pub favourites: FavouriteSync,
pub playlists: crate::playlists::PlaylistSync,
pub history: crate::remote::history::HistorySync,
pub dsp: crate::remote::dsp_sync::DspSync,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Walk {
Always,
IfChanged,
}
pub fn sync_remote(
db: &Database,
client: &SubsonicClient,
walk: Walk,
url: &str,
username: &str,
progress: &(dyn Fn(crate::remote::sync::SyncProgress) + Sync),
) -> Result<Synced, crate::remote::sync::SyncError> {
use crate::remote::sync;
static SYNCING: parking_lot::Mutex<()> = parking_lot::Mutex::new(());
let _one_at_a_time = SYNCING.lock();
let walked = sync::library_version(db, url);
let modified = client.library_modified(walked);
crate::remote::refusal::observe(client.auth(), &modified);
let version = modified
.inspect_err(|e| log::debug!("library version unavailable: {e}"))
.ok()
.flatten();
let library = if walk == Walk::IfChanged && version.is_some() && version == walked {
log::info!("library unchanged on the server; not walked");
sync::SyncResult::default()
} else {
let library = sync::sync_library(db, client, url, username, progress).inspect_err(|e| {
if let sync::SyncError::Subsonic(e) = e {
crate::remote::refusal::observe_error(client.auth(), e);
}
})?;
if library.is_complete()
&& let Some(version) = version
{
sync::set_library_version(db, url, username, version)?;
}
library
};
Ok(Synced {
library,
favourites: reconcile_favourites(db, client),
playlists: crate::playlists::reconcile_playlists(db, client, url, username),
history: crate::remote::history::reconcile(db, client, url, username),
dsp: crate::remote::dsp_sync::reconcile(db, client, url),
})
}
#[derive(Debug, Default, Clone, Copy)]
pub struct FavouriteSync {
pub pushed: usize,
pub imported: usize,
}
pub fn reconcile_favourites(db: &Database, client: &SubsonicClient) -> FavouriteSync {
let mut out = FavouriteSync::default();
for change in queries::favourite_changes(&db.conn).unwrap_or_default() {
let kind = match change.kind.as_str() {
"album" => FavouriteKind::Album,
"artist" => FavouriteKind::Artist,
_ => FavouriteKind::Track,
};
let sent = send_favourite(client, kind, &change.remote_id, change.star);
match &sent {
Ok(()) => out.pushed += 1,
Err(SubsonicError::Api { code: 40..=44, .. }) => {
log::warn!("favourite change refused for the account; kept for the next sync");
continue;
}
Err(e @ SubsonicError::Api { .. }) => {
log::warn!("favourite change refused by the server; dropped: {e}");
}
Err(e) => {
log::warn!("favourite change not sent; kept for the next sync: {e}");
continue;
}
}
if let Err(e) = queries::forget_favourite_change(&db.conn, &change) {
log::warn!("could not clear a sent favourite change: {e}");
}
}
let starred = match client.get_starred_all() {
Ok(s) => s,
Err(e) => {
log::warn!("could not fetch starred items from the server: {e}");
return out;
}
};
let songs: Vec<String> = starred.song.into_iter().map(|s| s.id).collect();
let albums: Vec<String> = starred.album.into_iter().map(|a| a.id).collect();
let artists: Vec<String> = starred.artist.into_iter().map(|a| a.id).collect();
let unstarred = |ids: Vec<String>, starred: &[String]| {
let starred: std::collections::HashSet<&String> = starred.iter().collect();
ids.into_iter()
.filter(|id| !starred.contains(id))
.collect::<Vec<_>>()
};
let tracks = queries::favourites_with_remote_id(&db.conn, queries::LOCAL_USER)
.unwrap_or_default()
.into_iter()
.map(|(_, id)| id)
.collect();
for remote_id in unstarred(tracks, &songs) {
if client.star(&remote_id).is_ok() {
out.pushed += 1;
}
}
let local_albums = queries::favourite_albums_with_remote_id(&db.conn, queries::LOCAL_USER)
.unwrap_or_default()
.into_iter()
.map(|(_, id)| id)
.collect();
for remote_id in unstarred(local_albums, &albums) {
if client.star_album(&remote_id).is_ok() {
out.pushed += 1;
}
}
let local_artists = queries::favourite_artists_with_remote_id(&db.conn, queries::LOCAL_USER)
.unwrap_or_default()
.into_iter()
.map(|(_, id)| id)
.collect();
for remote_id in unstarred(local_artists, &artists) {
if client.star_artist(&remote_id).is_ok() {
out.pushed += 1;
}
}
let imported = queries::atomically(&db.conn, || -> rusqlite::Result<usize> {
let changed: std::collections::HashSet<(String, String)> =
queries::favourite_changes(&db.conn)?
.into_iter()
.map(|c| (c.kind, c.remote_id))
.collect();
let settled = |kind: FavouriteKind, ids: &[String]| -> Vec<String> {
ids.iter()
.filter(|id| !changed.contains(&(kind.as_str().to_owned(), (*id).clone())))
.cloned()
.collect()
};
let user = queries::LOCAL_USER;
Ok(queries::import_remote_favourites(
&db.conn,
user,
&settled(FavouriteKind::Track, &songs),
)? + queries::import_remote_favourite_albums(
&db.conn,
user,
&settled(FavouriteKind::Album, &albums),
)? + queries::import_remote_favourite_artists(
&db.conn,
user,
&settled(FavouriteKind::Artist, &artists),
)?)
});
match imported {
Ok(n) => out.imported += n,
Err(e) => log::warn!("could not import the server's favourites: {e}"),
}
out
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FavouriteKind {
Track,
Album,
Artist,
}
impl FavouriteKind {
fn as_str(self) -> &'static str {
match self {
Self::Track => "track",
Self::Album => "album",
Self::Artist => "artist",
}
}
}
pub fn sync_collection_favourite_to_remote(
db: &Database,
kind: FavouriteKind,
id: i64,
star: bool,
) {
let cfg = Config::load().unwrap_or_default();
if !cfg.remote.enabled {
return;
}
let remote_id = match kind {
FavouriteKind::Album => queries::album_remote_id(&db.conn, id),
FavouriteKind::Artist => queries::artist_remote_id(&db.conn, id),
FavouriteKind::Track => return,
};
let Ok(Some(remote_id)) = remote_id else {
log::warn!("not syncing favourite: {kind:?} {id} has no remote id");
return;
};
push_favourite(db, &cfg, kind, remote_id, star);
}
#[derive(Debug, thiserror::Error)]
pub enum SignInError {
#[error("the server did not accept those credentials: {0}")]
Rejected(#[from] crate::remote::client::SubsonicError),
#[error(
"this server needs an app password or API key: make one in the server's web UI and sign in with it"
)]
NeedsKey,
#[error("could not write the configuration: {0}")]
Config(#[from] crate::config::ConfigError),
}
pub fn set_remote_credentials(
url: &str,
username: &str,
password: &str,
) -> Result<(), SignInError> {
use crate::remote::client::{koan_sign_in, offers_unsigned};
let url = url.trim_end_matches('/');
if offers_unsigned(url, crate::remote::profile::SIGN_IN).unwrap_or(false) {
let device = crate::remote::link::LinkIdentity::this_device(None).name;
match koan_sign_in(url, username, password, &device) {
Ok(joined) => return adopt_api_key(url, &joined.username, &joined.api_key),
Err(SubsonicError::Api { code: 50, .. }) => {}
Err(e) => return Err(rejected(e)),
}
}
SubsonicClient::new(url, username, password)
.ping()
.map_err(rejected)?;
remember_remote(url, username, Credential::Password(password.to_string()))
}
pub fn set_remote_api_key(url: &str, username: &str, api_key: &str) -> Result<(), SignInError> {
let url = url.trim_end_matches('/');
let credential = Credential::ApiKey(api_key.to_string());
SubsonicClient::from_auth(SubsonicAuth::with(url, username, credential.clone())).ping()?;
remember_remote(url, username, credential)
}
pub fn change_own_password(current: &str, password: &str) -> Result<(), SignInError> {
let cfg = Config::load()?;
let client = subsonic_client(&cfg).ok_or(SignInError::Rejected(SubsonicError::BadResponse))?;
let device = crate::remote::link::LinkIdentity::this_device(None).name;
let joined = client.koan_change_own_password(current, password, &device)?;
adopt_api_key(&cfg.remote.url, &joined.username, &joined.api_key)
}
fn rejected(e: SubsonicError) -> SignInError {
match e {
SubsonicError::Api { code: 41, .. } => SignInError::NeedsKey,
e => SignInError::Rejected(e),
}
}
pub fn join_with_invite(invite: &crate::invite::Invite) -> Result<(), SignInError> {
let url = invite.server.trim_end_matches('/');
match (&invite.token, &invite.password) {
(Some(token), _) => {
let device = crate::remote::link::LinkIdentity::this_device(None).name;
let joined = crate::remote::client::redeem_invite(url, token, &device)?;
adopt_api_key(url, &joined.username, &joined.api_key)
}
(None, Some(password)) => set_remote_credentials(url, &invite.username, password),
(None, None) => Err(SignInError::Rejected(SubsonicError::BadResponse)),
}
}
pub(crate) fn adopt_api_key(url: &str, username: &str, api_key: &str) -> Result<(), SignInError> {
let url = url.trim_end_matches('/');
let replaced = Config::load()
.ok()
.filter(|c| c.remote.url.trim_end_matches('/') == url)
.filter(|c| c.remote.username == username)
.map(|c| c.remote.api_key)
.filter(|k| !k.is_empty() && k != api_key);
let revoke = |key: &str| {
let credential = Credential::ApiKey(key.to_string());
let client = SubsonicClient::from_auth(SubsonicAuth::with(url, username, credential));
if let Err(e) = client.koan_revoke_own_key() {
log::warn!("could not revoke an unused API key: {e}");
}
};
if let Err(e) = remember_remote(url, username, Credential::ApiKey(api_key.to_string())) {
revoke(api_key);
return Err(e);
}
if let Some(old) = replaced {
revoke(&old);
}
Ok(())
}
fn remember_remote(url: &str, username: &str, credential: Credential) -> Result<(), SignInError> {
Config::persist(|cfg| {
cfg.remote.enabled = true;
cfg.remote.url = url.to_string();
cfg.remote.username = username.to_string();
(cfg.remote.password, cfg.remote.api_key) = match &credential {
Credential::Password(p) => (p.clone(), String::new()),
Credential::ApiKey(k) => (String::new(), k.clone()),
};
cfg.remote.device_key = match &credential {
Credential::ApiKey(_) => crate::remote::proof::new_device_key().unwrap_or_default(),
Credential::Password(_) => String::new(),
};
})?;
crate::remote::proof::forget();
crate::remote::link::relink();
crate::remote::nearby::readvertise();
Ok(())
}
pub fn get_subsonic_password(cfg: &Config) -> Option<String> {
(!cfg.subsonic.password.is_empty()).then(|| cfg.subsonic.password.clone())
}
pub fn subsonic_auth(cfg: &Config) -> Option<SubsonicAuth> {
if !cfg.remote.enabled || cfg.remote.url.is_empty() {
return None;
}
Some(SubsonicAuth::with(
&cfg.remote.url,
&cfg.remote.username,
remote_credential(cfg)?,
))
}
pub fn subsonic_client(cfg: &Config) -> Option<Arc<SubsonicClient>> {
let auth = subsonic_auth(cfg)?;
let mut slot = SUBSONIC_CLIENT.lock();
if let Some((cached, client)) = slot.as_ref()
&& *cached == auth
{
return Some(client.clone());
}
let client = Arc::new(SubsonicClient::from_auth(auth.clone()));
*slot = Some((auth, client.clone()));
Some(client)
}
type CachedClient = Option<(SubsonicAuth, Arc<SubsonicClient>)>;
static SUBSONIC_CLIENT: std::sync::LazyLock<parking_lot::Mutex<CachedClient>> =
std::sync::LazyLock::new(|| parking_lot::Mutex::new(None));
#[derive(Debug, thiserror::Error)]
pub enum ShareError {
#[error("sharing.public_url is not set, so there is no address to give out")]
NoPublicUrl,
#[error("none of these tracks are in the library")]
NothingToShare,
#[error("none of these tracks are on the server, so a link has nothing to point at")]
NothingRemote,
#[error("the server refused to share these: {0}")]
Server(#[from] crate::remote::client::SubsonicError),
#[error(transparent)]
Database(#[from] crate::db::connection::DbError),
}
#[derive(Debug, Clone)]
pub struct ShareOutcome {
pub url: String,
pub id: String,
pub shared: usize,
pub skipped: usize,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ShareTarget {
Tracks(Vec<i64>),
Album {
album_id: i64,
start_track_id: Option<i64>,
},
Artist(i64),
}
pub fn resolve_share(
conn: &rusqlite::Connection,
target: &ShareTarget,
) -> Result<(Slice, Vec<i64>), ShareError> {
let album_tracks = |album_id| -> Result<Vec<i64>, ShareError> {
Ok(queries::tracks_for_album(conn, album_id)?
.into_iter()
.map(|t| t.id)
.collect())
};
let (slice, ids) = match target {
ShareTarget::Tracks(ids) => {
let rows = queries::tracks_by_ids(conn, ids)?;
match (ids.as_slice(), rows.first().and_then(|t| t.album_id)) {
([one], Some(album_id)) => {
return resolve_share(
conn,
&ShareTarget::Album {
album_id,
start_track_id: Some(*one),
},
);
}
_ => {
let ids = ids
.iter()
.copied()
.filter(|id| rows.iter().any(|t| t.id == *id))
.collect();
(Slice::TRACKS, ids)
}
}
}
ShareTarget::Album {
album_id,
start_track_id,
} => {
let ids = album_tracks(*album_id)?;
let slice = Slice {
kind: ShareKind::Album,
subject_id: Some(*album_id),
start_track_id: start_track_id.filter(|s| ids.contains(s)),
};
(slice, ids)
}
ShareTarget::Artist(artist_id) => {
let mut ids = Vec::new();
for album in queries::albums_for_artist(conn, *artist_id)? {
ids.extend(album_tracks(album.id)?);
}
let slice = Slice {
kind: ShareKind::Artist,
subject_id: Some(*artist_id),
start_track_id: None,
};
(slice, ids)
}
};
if ids.is_empty() {
return Err(ShareError::NothingToShare);
}
Ok((slice, ids))
}
pub fn create_share(
db: &Database,
user: i64,
cfg: &Config,
target: &ShareTarget,
description: Option<&str>,
) -> Result<ShareOutcome, ShareError> {
let Some(client) = subsonic_client(cfg) else {
return create_native_share(db, user, cfg, target, description);
};
let resolved;
let track_ids = match target {
ShareTarget::Tracks(ids) => ids.as_slice(),
_ => {
resolved = resolve_share(&db.conn, target)?.1;
resolved.as_slice()
}
};
let rows = queries::tracks_by_ids(&db.conn, track_ids)?;
let shared = rows.iter().filter(|t| t.remote_id.is_some()).count();
if shared == 0 {
return Err(ShareError::NothingRemote);
}
let one_album = rows
.first()
.and_then(|f| f.album_id)
.filter(|first| rows.iter().all(|t| t.album_id == Some(*first)))
.and_then(|album_id| album_remote_id(&db.conn, album_id, rows.len()));
let remote_ids: Vec<String> = match one_album {
Some(rid) => vec![album_share_id(&client, rid)],
None => rows.into_iter().filter_map(|t| t.remote_id).collect(),
};
let refs: Vec<&str> = remote_ids.iter().map(String::as_str).collect();
let share = client.create_share(&refs, description)?;
let url = share
.url
.clone()
.unwrap_or_else(|| format!("{}/s/{}", client.base_url(), share.id));
Ok(ShareOutcome {
url,
id: share.id,
shared,
skipped: track_ids.len().saturating_sub(shared),
})
}
pub fn create_native_share(
db: &Database,
user: i64,
cfg: &Config,
target: &ShareTarget,
description: Option<&str>,
) -> Result<ShareOutcome, ShareError> {
let base = cfg
.sharing
.public_url
.as_deref()
.filter(|u| !u.trim().is_empty())
.ok_or(ShareError::NoPublicUrl)?;
let (slice, ids) = resolve_share(&db.conn, target)?;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_secs() as i64);
let share = queries::shares::create_share(&db.conn, user, slice, &ids, description, now, None)?;
Ok(ShareOutcome {
url: share_url(base, &share.id),
id: share.id,
shared: ids.len(),
skipped: match (target, slice.kind) {
(ShareTarget::Tracks(asked), ShareKind::Tracks) => asked.len() - ids.len(),
_ => 0,
},
})
}
pub fn share_url(public_url: &str, id: &str) -> String {
format!("{}/share/{id}", public_url.trim_end_matches('/'))
}
fn album_share_id(client: &crate::remote::client::SubsonicClient, remote_id: String) -> String {
let koan = crate::remote::profile::is_koan(client.auth());
album_share_id_for(koan, remote_id)
}
fn album_share_id_for(koan: bool, remote_id: String) -> String {
if koan && remote_id.parse::<i64>().is_ok() {
format!("al-{remote_id}")
} else {
remote_id
}
}
fn album_remote_id(conn: &rusqlite::Connection, album_id: i64, selected: usize) -> Option<String> {
let (remote_id, total): (Option<String>, i64) = conn
.query_row(
"SELECT al.remote_id, (SELECT COUNT(*) FROM tracks WHERE album_id = al.id)
FROM albums al WHERE al.id = ?1",
[album_id],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.ok()?;
(total == selected as i64).then_some(remote_id).flatten()
}
pub fn shuffle<T>(items: &mut [T]) {
let mut seed = [0u8; 8];
if getrandom::fill(&mut seed).is_err() {
return; }
let mut state = u64::from_le_bytes(seed) | 1;
for i in (1..items.len()).rev() {
state ^= state << 13;
state ^= state >> 7;
state ^= state << 17;
items.swap(i, (state % (i as u64 + 1)) as usize);
}
}
pub fn truncate_bytes(s: &str, max: usize) -> &str {
if s.len() <= max {
return s;
}
let mut end = max;
while end > 0 && !s.is_char_boundary(end) {
end -= 1;
}
&s[..end]
}
pub fn sanitise_filename(s: &str) -> String {
let cleaned: String = s
.chars()
.map(|c| match c {
'/' | '\\' | ':' | '*' | '?' | '"' | '<' | '>' | '|' => '_',
_ => c,
})
.collect::<String>()
.trim()
.to_string();
let cleaned = truncate_bytes(&cleaned, 240).trim_end().to_string();
match cleaned.as_str() {
"." | ".." => "_".into(),
_ => cleaned,
}
}
pub fn sanitise_extension(codec: &str) -> Option<String> {
let ext: String = codec
.chars()
.filter(char::is_ascii_alphanumeric)
.take(16)
.collect::<String>()
.to_lowercase();
(!ext.is_empty()).then_some(ext)
}
pub fn path_within(dir: &Path, path: &Path) -> bool {
path.strip_prefix(dir).is_ok_and(|rest| {
rest.components()
.all(|c| matches!(c, std::path::Component::Normal(_)))
})
}
pub fn year_of(date: &str) -> Option<&str> {
date.get(..4)
}
pub fn cache_path_for_track(
cache_dir: &Path,
track: &queries::TrackRow,
album_date: Option<&str>,
) -> PathBuf {
let artist_dir = sanitise_filename(&track.artist_name);
let year = album_date
.and_then(year_of)
.map(|y| format!("({}) ", y))
.unwrap_or_default();
let codec = track
.codec
.as_deref()
.map(|c| format!(" [{}]", c))
.unwrap_or_default();
let album_dir = sanitise_filename(&format!("{}{}{}", year, track.album_title, codec));
let disc_prefix = match track.disc {
Some(d) if d > 1 => format!("{}-", d),
_ => String::new(),
};
let track_num = track
.track_number
.map(|n| format!("{:02}. ", n))
.unwrap_or_default();
let ext = track
.codec
.as_deref()
.and_then(sanitise_extension)
.unwrap_or_else(|| "flac".into());
let filename = sanitise_filename(&format!(
"{}{}{} - {}",
disc_prefix, track_num, track.artist_name, track.title
));
cache_dir
.join(artist_dir)
.join(album_dir)
.join(format!("{}.{}", filename, ext))
}
fn resolve_item_path(
cfg: &Config,
track: &queries::TrackRow,
remote_url: Option<&str>,
album_date: Option<&str>,
) -> (PathBuf, ItemState) {
match queries::choose_playback_source(
track.path.as_deref(),
track.cached_path.as_deref(),
remote_url,
) {
Some(queries::PlaybackSource::Local(p)) => (p, ItemState::Ready),
Some(queries::PlaybackSource::Cached(p)) => {
let state = if is_cached_audio(&p) {
ItemState::Ready
} else {
ItemState::Pending
};
(p, state)
}
Some(queries::PlaybackSource::Remote(_)) => {
let dest = cache_path_for_track(&cfg.cache_dir(), track, album_date);
if dest.exists() && is_cached_audio(&dest) {
(dest, ItemState::Ready)
} else {
(dest, ItemState::Pending)
}
}
_ => {
let dest = cache_path_for_track(&cfg.cache_dir(), track, album_date);
(dest, ItemState::Pending)
}
}
}
pub fn playlist_item_from_track(
track: &queries::TrackRow,
album_date: Option<&str>,
dest: PathBuf,
state: ItemState,
) -> PlaylistItem {
let year = album_date.and_then(year_of).map(str::to_string);
PlaylistItem {
playlist_entry_id: None,
id: QueueItemId::new(),
db_id: Some(track.id),
path: dest,
title: track.title.clone(),
artist: track.artist_name.clone(),
album_artist: track.album_artist_name.clone(),
album: track.album_title.clone(),
year,
codec: track.codec.clone(),
track_number: track.track_number.map(|n| n as i64),
disc: track.disc.map(|n| n as i64),
duration_ms: track.duration_ms.map(|d| d as u64),
state,
pre_shuffle: None,
}
}
pub fn playlist_items_for_tracks(db: &Database, tracks: &[queries::TrackRow]) -> Vec<PlaylistItem> {
let cfg = Config::load().unwrap_or_default();
let ids: Vec<i64> = tracks.iter().map(|t| t.id).collect();
let extras = queries::queue_item_extras(&db.conn, &ids).unwrap_or_default();
tracks
.iter()
.map(|track| {
let extra = extras.get(&track.id);
let remote_url = extra.and_then(|e| e.remote_url.as_deref());
let album_date = extra.and_then(|e| e.album_date.as_deref());
let (path, state) = resolve_item_path(&cfg, track, remote_url, album_date);
playlist_item_from_track(track, album_date, path, state)
})
.collect()
}
fn is_cached_audio(path: &std::path::Path) -> bool {
const MIN_PLAUSIBLE_BYTES: u64 = 4096;
match std::fs::metadata(path) {
Ok(meta) if meta.len() >= MIN_PLAUSIBLE_BYTES => true,
Ok(_) => {
let mut first = [0u8; 1];
match std::fs::File::open(path)
.and_then(|mut f| std::io::Read::read_exact(&mut f, &mut first).map(|_| first[0]))
{
Ok(b) => b != b'{' && b != b'<',
Err(_) => false,
}
}
Err(_) => false,
}
}
pub(crate) fn download_track(
db_id: i64,
cancelled: &dyn Fn() -> bool,
tx: &crossbeam_channel::Sender<PlayerCommand>,
state: &SharedPlayerState,
cfg: &Config,
client: &SubsonicClient,
) -> Option<Result<PathBuf, String>> {
let db = match crate::db::pool::shared().get() {
Ok(db) => db,
Err(e) => return Some(Err(format!("db error: {e}"))),
};
let Ok(Some(track)) = queries::get_track_row(&db.conn, db_id) else {
return Some(Err("track not found".into()));
};
if let Some(p) = track.path.as_deref().map(PathBuf::from)
&& p.exists()
{
log::info!("download_track: local file exists, using {}", p.display());
return Some(Ok(p));
}
let Some(remote_id) = track.remote_id.clone() else {
return Some(Err(
"not in the library folder, and no remote copy to fetch".into(),
));
};
let album_date: Option<String> = track
.album_id
.and_then(|aid| queries::album_date(&db.conn, aid).ok().flatten());
let cache_dir = cfg.cache_dir();
let dest = cache_path_for_track(&cache_dir, &track, album_date.as_deref());
if !path_within(&cache_dir, &dest) {
return Some(Err(format!(
"cache path escapes the cache: {}",
dest.display()
)));
}
if dest.exists() && !is_cached_audio(&dest) {
log::warn!(
"discarding non-audio cache entry {} (likely a stored server error)",
dest.display()
);
let _ = std::fs::remove_file(&dest);
}
if dest.exists() {
return Some(Ok(dest));
}
let store = state.downloads();
let bytes_written = store.announce(
db_id,
track.title.clone(),
track.artist_name.clone(),
crate::remote::download::part_path(&dest),
dest.clone(),
);
let progress_tx = tx.clone();
let stream_ready_flag = std::sync::atomic::AtomicBool::new(false);
let announced_total = AtomicU64::new(u64::MAX);
let result =
client.download_with_progress(&remote_id, &dest, cancelled, |downloaded, total| {
bytes_written.set(downloaded);
store.progressed();
if announced_total.swap(total, Ordering::Relaxed) != total {
store.started(db_id, total);
}
if !stream_ready_flag.load(Ordering::Relaxed)
&& downloaded >= crate::player::state::STREAM_THRESHOLD
{
stream_ready_flag.store(true, Ordering::Relaxed);
for id in store.waiters(db_id) {
progress_tx.send(PlayerCommand::TrackStreamReady(id)).ok();
}
}
});
match result {
Err(SubsonicError::Download(DownloadError::Cancelled)) => None,
Err(e) => {
log::warn!("x {} — {}", track.title, e);
Some(Err(e.to_string()))
}
Ok(()) => {
cache_grew(std::fs::metadata(&dest).map_or(0, |m| m.len()));
if let Err(e) = queries::set_cached_path(&db.conn, db_id, &dest.to_string_lossy()) {
log::warn!(
"cached {} but failed to record it ({}) — it will not be evicted",
dest.display(),
e
);
}
log::info!("+ {} — {}", track.title, track.artist_name);
Some(Ok(dest))
}
}
}
pub fn remote_unavailable(cfg: &Config) -> String {
if !cfg.remote.enabled {
return "no remote server is configured".into();
}
if cfg.remote.url.is_empty() {
return "the remote server has no address".into();
}
if remote_credential(cfg).is_none() {
return "no password or API key is stored for the remote server".into();
}
"the remote server could not be reached".into()
}
pub const SIGN_IN_REFUSED: &str = "the remote server refused the stored sign-in; sign in again";
pub fn remote_problem(cfg: &Config) -> Option<String> {
if !cfg.remote.enabled || cfg.remote.url.is_empty() {
return None;
}
match subsonic_auth(cfg) {
None => Some(remote_unavailable(cfg)),
Some(auth) if crate::remote::refusal::refused(&auth) => Some(SIGN_IN_REFUSED.into()),
Some(_) => None,
}
}
pub fn sign_in_refused(cfg: &Config) -> bool {
subsonic_auth(cfg).is_some_and(|auth| crate::remote::refusal::refused(&auth))
}
#[cfg(test)]
mod year_tests {
use super::year_of;
#[test]
fn a_year_is_the_first_four_characters_when_they_are_bytes_too() {
assert_eq!(year_of("1997-05-21"), Some("1997"));
assert_eq!(year_of("199"), None);
assert_eq!(year_of("1997"), None);
}
}
#[cfg(test)]
mod rebuild_tests {
use super::*;
use crate::db::queries::sample_meta;
fn test_db() -> Database {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.pragma_update(None, "foreign_keys", "on").unwrap();
crate::db::schema::create_tables(&conn).unwrap();
Database { conn }
}
#[test]
fn cached_paths_follow_a_moved_cache_directory() {
let old = tempfile::tempdir().unwrap();
let new = tempfile::tempdir().unwrap();
let db = test_db();
let mut rows = Vec::new();
for name in ["moved", "gone", "current"] {
let mut meta = sample_meta(name, "Artist", "Album");
meta.source = "remote".into();
meta.path = None;
meta.remote_id = Some(name.into());
let id = queries::upsert_track(&db.conn, &meta).unwrap();
let tail = format!("Artist/Album/{name}.flac");
if name != "gone" {
let file = new.path().join(&tail);
std::fs::create_dir_all(file.parent().unwrap()).unwrap();
std::fs::write(&file, b"audio").unwrap();
}
let stored = if name == "current" {
new.path()
} else {
old.path()
}
.join(&tail);
queries::set_cached_path(&db.conn, id, &stored.to_string_lossy()).unwrap();
rows.push((id, tail));
}
assert_eq!(relocate_cached_paths(&db, new.path()).unwrap(), 1);
let cached = |id: i64| -> String {
db.conn
.query_row("SELECT cached_path FROM tracks WHERE id = ?1", [id], |r| {
r.get(0)
})
.unwrap()
};
let expect = |root: &Path, tail: &str| root.join(tail).to_string_lossy().into_owned();
assert_eq!(
cached(rows[0].0),
expect(new.path(), &rows[0].1),
"re-rooted"
);
assert_eq!(
cached(rows[1].0),
expect(old.path(), &rows[1].1),
"no file, left alone"
);
assert_eq!(
cached(rows[2].0),
expect(new.path(), &rows[2].1),
"already current"
);
assert_eq!(
relocate_cached_paths(&db, new.path()).unwrap(),
0,
"idempotent"
);
}
#[test]
fn clearing_one_download_leaves_the_others_and_the_library_alone() {
let dir = tempfile::tempdir().unwrap();
let db = test_db();
let mut cached = Vec::new();
for name in ["one", "two"] {
let mut meta = sample_meta(name, "Artist", "Album");
meta.source = "remote".into();
meta.path = None;
meta.remote_id = Some(name.into());
let id = queries::upsert_track(&db.conn, &meta).unwrap();
let file = dir.path().join(format!("{name}.opus"));
std::fs::write(&file, vec![0u8; 2048]).unwrap();
queries::set_cached_path(&db.conn, id, &file.to_string_lossy()).unwrap();
cached.push((id, file));
}
let cleared = clear_downloads_for(&db, &[cached[0].0]);
assert_eq!(cleared.files, 1);
assert_eq!(cleared.bytes, 2048);
assert!(!cached[0].1.exists(), "the copy asked for is gone");
assert!(cached[1].1.exists(), "the other one is untouched");
assert_eq!(queries::library_stats(&db.conn).unwrap().remote_tracks, 2);
assert_eq!(queries::library_stats(&db.conn).unwrap().cached_tracks, 1);
assert!(
queries::cached_paths_for(&db.conn, &[cached[0].0])
.unwrap()
.is_empty()
);
}
#[test]
fn clearing_a_download_that_is_already_gone_is_not_a_failure() {
let db = test_db();
let mut meta = sample_meta("ghost", "Artist", "Album");
meta.source = "remote".into();
meta.path = None;
meta.remote_id = Some("ghost".into());
let id = queries::upsert_track(&db.conn, &meta).unwrap();
queries::set_cached_path(&db.conn, id, "/nowhere/at/all.opus").unwrap();
let cleared = clear_downloads_for(&db, &[id]);
assert_eq!(cleared.files, 0, "nothing was there to remove");
assert!(
queries::cached_paths_for(&db.conn, &[id])
.unwrap()
.is_empty()
);
}
#[test]
fn sweeping_removes_half_finished_downloads_and_nothing_else() {
let dir = tempfile::tempdir().unwrap();
let cache = dir.path().join("cache");
std::fs::create_dir_all(cache.join("Artist")).unwrap();
let finished = cache.join("Artist/whole.opus");
let half = cache.join("Artist/half.opus.part");
std::fs::write(&finished, vec![0u8; 1024]).unwrap();
std::fs::write(&half, vec![0u8; 4096]).unwrap();
let cfg = Config {
remote: crate::config::RemoteConfig {
cache_dir: Some(cache.clone()),
..Default::default()
},
..Default::default()
};
let swept = sweep_partial_downloads(&cfg);
assert_eq!(swept.files, 1);
assert_eq!(swept.bytes, 4096);
assert!(!half.exists(), "the unfinished one is gone");
assert!(finished.exists(), "a downloaded track is not touched");
}
#[test]
fn the_cache_is_measured_as_a_clear_counts_it() {
let dir = tempfile::tempdir().unwrap();
let cache = dir.path().join("cache");
std::fs::create_dir_all(cache.join("Artist/Album")).unwrap();
std::fs::write(cache.join("Artist/Album/01.flac"), vec![0u8; 3000]).unwrap();
std::fs::write(cache.join("Artist/Album/02.flac"), vec![0u8; 5000]).unwrap();
let cfg = Config {
remote: crate::config::RemoteConfig {
cache_dir: Some(cache.clone()),
..Default::default()
},
..Default::default()
};
let db = Database::open(&dir.path().join("koan.db")).unwrap();
assert_eq!(measure_cache(&cfg), 8000);
let cleared = clear_download_cache(&db, &cfg);
assert_eq!((cleared.files, cleared.bytes), (2, 8000));
assert_eq!(measure_cache(&cfg), 0, "nothing left to count");
}
#[test]
fn sweeping_an_empty_cache_is_not_an_error() {
let dir = tempfile::tempdir().unwrap();
let cfg = Config {
remote: crate::config::RemoteConfig {
cache_dir: Some(dir.path().join("nothing-here")),
..Default::default()
},
..Default::default()
};
assert_eq!(sweep_partial_downloads(&cfg).files, 0);
}
#[test]
fn clearing_no_tracks_does_nothing() {
let db = test_db();
assert_eq!(clear_downloads_for(&db, &[]).files, 0);
}
#[test]
fn a_rebuild_re_reads_the_library_into_the_rows_it_has() {
let db = test_db();
let mut meta = sample_meta("Windowlicker", "Aphex Twin", "Windowlicker EP");
meta.path = Some("/music/windowlicker.flac".into());
let track_id = queries::upsert_track(&db.conn, &meta).unwrap();
queries::toggle_favourite(&db.conn, crate::db::queries::LOCAL_USER, track_id).unwrap();
db.conn
.execute(
"INSERT INTO lyrics_cache (track_id, source, content, fetched_at)
VALUES (?1, 'test', 'la la la', 0)",
[track_id],
)
.unwrap();
let summary = rebuild_index(&db).unwrap();
assert_eq!(summary.tracks, 1);
assert_eq!(summary.albums, 1);
let count = |sql: &str| -> i64 { db.conn.query_row(sql, [], |r| r.get(0)).unwrap() };
assert_eq!(
count("SELECT COUNT(*) FROM local_files"),
0,
"every file is read again"
);
assert_eq!(queries::upsert_track(&db.conn, &meta).unwrap(), track_id);
assert_eq!(count("SELECT COUNT(*) FROM tracks"), 1);
assert_eq!(count("SELECT COUNT(*) FROM favourites"), 1);
assert_eq!(count("SELECT COUNT(*) FROM lyrics_cache"), 1);
}
#[test]
fn a_rebuilt_file_that_is_gone_goes_with_its_folder_scan() {
let db = test_db();
let tmp = tempfile::tempdir().unwrap();
let mut meta = sample_meta("Windowlicker", "Aphex Twin", "Windowlicker");
meta.path = Some(tmp.path().join("gone.flac").to_string_lossy().into_owned());
queries::upsert_track(&db.conn, &meta).unwrap();
rebuild_index(&db).unwrap();
queries::remove_stale_tracks(&db.conn, tmp.path(), false).unwrap();
let tracks: i64 = db
.conn
.query_row("SELECT COUNT(*) FROM tracks", [], |r| r.get(0))
.unwrap();
assert_eq!(tracks, 0);
}
#[test]
fn rebuilding_an_empty_library_is_not_an_error() {
let db = test_db();
let summary = rebuild_index(&db).unwrap();
assert_eq!(summary.tracks, 0);
}
}
#[cfg(test)]
mod share_tests {
use super::*;
use crate::db::queries::sample_meta;
fn test_db() -> Database {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.pragma_update(None, "foreign_keys", "on").unwrap();
crate::db::schema::create_tables(&conn).unwrap();
Database { conn }
}
fn album_of_three(db: &Database) -> (i64, Vec<i64>) {
let ids: Vec<i64> = ["One", "Two", "Three"]
.iter()
.enumerate()
.map(|(i, title)| {
let mut meta = sample_meta(title, "Boards of Canada", "Geogaddi");
meta.path = Some(format!("/music/geogaddi/{i}.flac"));
meta.track_number = Some(i as i32 + 1);
queries::upsert_track(&db.conn, &meta).unwrap()
})
.collect();
let album_id: i64 = db
.conn
.query_row("SELECT album_id FROM tracks WHERE id = ?1", [ids[0]], |r| {
r.get(0)
})
.unwrap();
db.conn
.execute(
"UPDATE albums SET remote_id = 'al-1' WHERE id = ?1",
[album_id],
)
.unwrap();
(album_id, ids)
}
#[test]
fn whole_album_collapses_to_the_album_link() {
let db = test_db();
let (album_id, ids) = album_of_three(&db);
assert_eq!(
album_remote_id(&db.conn, album_id, ids.len()),
Some("al-1".into())
);
}
#[test]
fn part_of_an_album_does_not() {
let db = test_db();
let (album_id, _) = album_of_three(&db);
assert_eq!(album_remote_id(&db.conn, album_id, 2), None);
}
#[test]
fn a_local_only_album_has_no_link_to_collapse_to() {
let db = test_db();
let (album_id, ids) = album_of_three(&db);
db.conn
.execute(
"UPDATE albums SET remote_id = NULL WHERE id = ?1",
[album_id],
)
.unwrap();
assert_eq!(album_remote_id(&db.conn, album_id, ids.len()), None);
}
}
#[cfg(test)]
mod client_cache_tests {
use super::*;
#[test]
fn one_subsonic_client_is_shared_per_credentials() {
crate::config::isolate_config_for_tests();
let mut cfg = Config::default();
cfg.remote.enabled = true;
cfg.remote.url = "https://shared-client.invalid".into();
cfg.remote.username = "koan".into();
cfg.remote.password = "first".into();
let first = subsonic_client(&cfg).expect("a configured remote yields a client");
let again = subsonic_client(&cfg).expect("a configured remote yields a client");
assert!(
Arc::ptr_eq(&first, &again),
"rebuilding drops the connection pool and re-handshakes TLS per request"
);
cfg.remote.password = "second".into();
let relogged = subsonic_client(&cfg).expect("a configured remote yields a client");
assert!(
!Arc::ptr_eq(&first, &relogged),
"new credentials must not keep serving the client signed with the old ones"
);
}
}
#[cfg(test)]
mod native_share_tests {
use super::*;
use crate::db::queries::{sample_meta, upsert_track};
#[test]
fn a_standalone_server_shares_natively_in_the_order_asked() {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.pragma_update(None, "foreign_keys", "on").unwrap();
crate::db::schema::create_tables(&conn).unwrap();
let db = Database { conn };
let a = upsert_track(&db.conn, &sample_meta("A", "X", "Y")).unwrap();
let b = upsert_track(&db.conn, &sample_meta("B", "X", "Y")).unwrap();
let mut cfg = Config::default();
assert!(matches!(
create_share(
&db,
queries::LOCAL_USER,
&cfg,
&ShareTarget::Tracks(vec![a]),
None
),
Err(ShareError::NoPublicUrl)
));
cfg.sharing.public_url = Some("https://koan.example/".into());
let out = create_share(
&db,
queries::LOCAL_USER,
&cfg,
&ShareTarget::Tracks(vec![b, 9999, a]),
Some("mix"),
)
.unwrap();
assert_eq!(out.url, format!("https://koan.example/share/{}", out.id));
assert_eq!((out.shared, out.skipped), (2, 1));
let share = queries::shares::get_share(&db.conn, &out.id)
.unwrap()
.unwrap();
assert_eq!(share.track_ids, [b, a]);
assert!(matches!(
create_share(
&db,
queries::LOCAL_USER,
&cfg,
&ShareTarget::Tracks(vec![9999]),
None
),
Err(ShareError::NothingToShare)
));
}
#[test]
fn a_server_with_an_upstream_still_shares_natively() {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.pragma_update(None, "foreign_keys", "on").unwrap();
crate::db::schema::create_tables(&conn).unwrap();
let db = Database { conn };
let a = upsert_track(&db.conn, &sample_meta("A", "X", "Y")).unwrap();
let mut cfg = Config::default();
cfg.remote.enabled = true;
cfg.remote.url = "https://upstream.invalid".into();
cfg.remote.username = "someone".into();
cfg.remote.password = "secret".into();
cfg.sharing.public_url = Some("https://koan.example".into());
let out = create_native_share(
&db,
queries::LOCAL_USER,
&cfg,
&ShareTarget::Tracks(vec![a]),
None,
)
.unwrap();
assert_eq!(out.url, format!("https://koan.example/share/{}", out.id));
assert!(
queries::shares::get_share(&db.conn, &out.id)
.unwrap()
.is_some()
);
}
fn album_track(db: &Database, title: &str, album: &str, n: i32, date: &str) -> i64 {
let mut meta = sample_meta(title, "Rrose", album);
meta.track_number = Some(n);
meta.date = Some(date.into());
upsert_track(&db.conn, &meta).unwrap()
}
#[test]
fn shares_are_slices_fixed_when_made() {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.pragma_update(None, "foreign_keys", "on").unwrap();
crate::db::schema::create_tables(&conn).unwrap();
let db = Database { conn };
let later = album_track(&db, "L1", "Later", 1, "2021");
let a1 = album_track(&db, "E1", "Earlier", 1, "2015");
let a2 = album_track(&db, "E2", "Earlier", 2, "2015");
let album_of = |t| {
queries::tracks_by_ids(&db.conn, &[t]).unwrap()[0]
.album_id
.unwrap()
};
let (earlier, later_album) = (album_of(a1), album_of(later));
let artist = queries::tracks_by_ids(&db.conn, &[a1]).unwrap()[0]
.artist_id
.unwrap();
let (slice, ids) = resolve_share(&db.conn, &ShareTarget::Tracks(vec![a2])).unwrap();
assert_eq!(
(slice.kind, slice.subject_id, slice.start_track_id),
(ShareKind::Album, Some(earlier), Some(a2))
);
assert_eq!(ids, [a1, a2]);
let (slice, ids) = resolve_share(
&db.conn,
&ShareTarget::Album {
album_id: later_album,
start_track_id: Some(a1),
},
)
.unwrap();
assert_eq!((slice.kind, slice.start_track_id), (ShareKind::Album, None));
assert_eq!(ids, [later]);
let (slice, ids) = resolve_share(&db.conn, &ShareTarget::Artist(artist)).unwrap();
assert_eq!(
(slice.kind, slice.subject_id),
(ShareKind::Artist, Some(artist))
);
assert_eq!(ids, [a1, a2, later]);
let (slice, ids) = resolve_share(&db.conn, &ShareTarget::Tracks(vec![later, a1])).unwrap();
assert_eq!(slice, Slice::TRACKS);
assert_eq!(ids, [later, a1]);
assert!(matches!(
resolve_share(&db.conn, &ShareTarget::Artist(9999)),
Err(ShareError::NothingToShare)
));
}
}
#[cfg(test)]
mod album_share_id_tests {
use super::album_share_id_for;
#[test]
fn a_koan_album_is_named_as_an_album() {
assert_eq!(album_share_id_for(true, "46215".into()), "al-46215");
assert_eq!(album_share_id_for(true, "al-7".into()), "al-7");
assert_eq!(album_share_id_for(false, "46215".into()), "46215");
assert_eq!(album_share_id_for(false, "3xJ9kQ2pZ".into()), "3xJ9kQ2pZ");
}
}
#[cfg(test)]
mod cache_path_tests {
use super::*;
fn track(artist: &str, album: &str, codec: &str) -> queries::TrackRow {
queries::TrackRow {
id: 1,
album_id: None,
artist_id: None,
artist_name: artist.into(),
album_artist_name: artist.into(),
album_title: album.into(),
disc: None,
track_number: Some(1),
title: "Song".into(),
duration_ms: None,
path: None,
codec: Some(codec.into()),
sample_rate: None,
bit_depth: None,
channels: None,
bitrate: None,
genre: None,
source: "remote".into(),
remote_id: Some("r1".into()),
cached_path: None,
}
}
#[test]
fn a_server_suffix_cannot_leave_the_cache() {
let cache = Path::new("/cache");
for codec in [
"flac/../../../../x",
"..",
"../..",
"/etc/passwd",
"\\..\\..",
] {
let path = cache_path_for_track(cache, &track("A", "B", codec), None);
assert!(path_within(cache, &path), "{codec}: {}", path.display());
}
let path = cache_path_for_track(cache, &track("A", "B", "flac/../../../../x"), None);
assert_eq!(path.extension().unwrap(), "flacx");
}
#[test]
fn dot_names_cannot_climb_out() {
let cache = Path::new("/cache");
let path = cache_path_for_track(cache, &track("..", ".", ".."), None);
assert!(path_within(cache, &path), "{}", path.display());
assert_eq!(sanitise_filename(".."), "_");
assert_eq!(sanitise_filename(" . "), "_");
assert_eq!(sanitise_filename("..."), "...");
}
#[test]
fn an_empty_suffix_falls_back_to_flac() {
assert_eq!(sanitise_extension("../"), None);
assert_eq!(sanitise_extension("FLAC"), Some("flac".into()));
let path = cache_path_for_track(Path::new("/c"), &track("A", "B", "./"), None);
assert_eq!(path.extension().unwrap(), "flac");
}
#[test]
fn path_within_rejects_parent_components() {
let dir = Path::new("/cache");
assert!(path_within(dir, Path::new("/cache/a/b.flac")));
assert!(!path_within(dir, Path::new("/cache/a/../../x")));
assert!(!path_within(dir, Path::new("/elsewhere/x")));
}
}
#[cfg(test)]
mod sign_in_tests {
use super::*;
fn serve() -> String {
use std::io::{BufRead, Write};
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
std::thread::spawn(move || {
for mut stream in listener.incoming().flatten() {
let mut reader = std::io::BufReader::new(stream.try_clone().unwrap());
let mut request = String::new();
reader.read_line(&mut request).unwrap();
let mut line = String::new();
while reader.read_line(&mut line).unwrap_or(0) > 2 {
line.clear();
}
let target = request.split_whitespace().nth(1).unwrap_or("");
let (path, query) = target.split_once('?').unwrap_or((target, ""));
let param = |name: &str| {
query
.split('&')
.filter_map(|kv| kv.split_once('='))
.find(|(k, _)| *k == name)
.map(|(_, v)| v.replace("%3A", ":"))
.unwrap_or_default()
};
let secret = match param("u").as_str() {
"mate" => "app-secret",
"testuser" => "shared-secret",
_ => "",
};
let ok = r#"{"subsonic-response":{"status":"ok"}}"#.to_owned();
let refused = |code: i32| {
format!(
r#"{{"subsonic-response":{{"status":"failed","error":{{"code":{code},"message":"refused"}}}}}}"#
)
};
let body = match path.rsplit('/').next().unwrap() {
"getOpenSubsonicExtensions" => r#"{"subsonic-response":{"status":"ok","openSubsonicExtensions":[{"name":"koanSignIn","versions":[1]}]}}"#.to_owned(),
"koanSignIn" => {
let hex = param("p").trim_start_matches("enc:").to_owned();
let typed: String = (0..hex.len())
.step_by(2)
.filter_map(|i| u8::from_str_radix(&hex[i..i + 2], 16).ok())
.map(char::from)
.collect();
match (param("u").as_str(), typed.as_str()) {
("mate", "hunter22") => r#"{"subsonic-response":{"status":"ok","join":{"username":"mate","apiKey":"minted"}}}"#.to_owned(),
(_, typed) if typed == secret => refused(50),
_ => refused(40),
}
}
"ping" if param("apiKey") == "minted" => ok,
"ping" => {
let expected =
format!("{:x}", md5::compute(format!("{secret}{}", param("s"))));
if !secret.is_empty() && param("t") == expected {
ok
} else {
refused(40)
}
}
_ => ok,
};
let _ = write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nConnection: close\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
}
});
url
}
#[test]
fn a_password_ends_in_a_key_and_other_credentials_are_kept_as_typed() {
let _guard = crate::config::tests::PERSIST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().unwrap();
crate::config::set_config_dir(dir.path());
let url = serve();
let kept = || {
let remote = Config::load().unwrap().remote;
(remote.username, remote.password, remote.api_key)
};
set_remote_credentials(&url, "mate", "hunter22").unwrap();
assert_eq!(kept(), ("mate".into(), String::new(), "minted".into()));
set_remote_credentials(&url, "mate", "app-secret").unwrap();
assert_eq!(
kept(),
("mate".into(), "app-secret".into(), String::new()),
"an app password is refused a key and kept"
);
set_remote_credentials(&url, "testuser", "shared-secret").unwrap();
assert_eq!(
kept(),
("testuser".into(), "shared-secret".into(), String::new()),
"the shared secret is refused a key and kept"
);
assert!(matches!(
set_remote_credentials(&url, "mate", "wrong"),
Err(SignInError::Rejected(SubsonicError::Api { code: 40, .. }))
));
assert_eq!(
kept().1,
"shared-secret",
"a refused sign-in writes nothing"
);
}
}
#[cfg(test)]
mod favourite_sync_tests {
use super::*;
use crate::db::queries::sample_meta;
use std::sync::Mutex;
#[derive(Clone, Copy)]
enum Unstar {
Taken,
Unreached,
Refused,
}
fn serve(stars: Arc<Mutex<Vec<String>>>) -> String {
serve_with(stars, Unstar::Taken)
}
fn serve_with(stars: Arc<Mutex<Vec<String>>>, unstar: Unstar) -> String {
use std::io::{BufRead, Write};
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
std::thread::spawn(move || {
for mut stream in listener.incoming().flatten() {
let mut reader = std::io::BufReader::new(stream.try_clone().unwrap());
let mut request = String::new();
reader.read_line(&mut request).unwrap();
let mut line = String::new();
while reader.read_line(&mut line).unwrap_or(0) > 2 {
line.clear();
}
let target = request.split_whitespace().nth(1).unwrap_or("");
let (path, query) = target.split_once('?').unwrap_or((target, ""));
let body = match path.rsplit('/').next().unwrap() {
"getStarred2" => {
r#"{"subsonic-response":{"status":"ok","starred2":{"song":[{"id":"s1","title":"One"}],"album":[{"id":"a1","name":"Album"}]}}}"#
}
"unstar" => match unstar {
Unstar::Taken => r#"{"subsonic-response":{"status":"ok"}}"#,
Unstar::Unreached => continue,
Unstar::Refused => {
r#"{"subsonic-response":{"status":"failed","error":{"code":70,"message":"not found"}}}"#
}
},
"star" => {
if let Some((_, id)) = query
.split('&')
.filter_map(|kv| kv.split_once('='))
.find(|(k, _)| *k == "id")
{
stars.lock().unwrap().push(id.to_string());
}
r#"{"subsonic-response":{"status":"ok"}}"#
}
_ => r#"{"subsonic-response":{"status":"ok"}}"#,
};
let _ = write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nConnection: close\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
}
});
url
}
#[test]
fn only_favourites_the_server_lacks_are_starred() {
let dir = tempfile::tempdir().unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
for (title, remote_id) in [("One", "s1"), ("Two", "s2")] {
let mut meta = sample_meta(title, "Artist", "Album");
meta.path = Some(format!("/music/{title}.flac"));
meta.remote_id = Some(remote_id.into());
let id = queries::upsert_track(&db.conn, &meta).unwrap();
queries::add_favourite(&db.conn, queries::LOCAL_USER, id).unwrap();
}
let stars = Arc::new(Mutex::new(Vec::new()));
let url = serve(stars.clone());
let sync = reconcile_favourites(&db, &SubsonicClient::new(&url, "u", "pw"));
assert_eq!(sync.pushed, 1);
assert_eq!(*stars.lock().unwrap(), ["s2"]);
}
fn album(db: &Database) -> i64 {
let mut meta = sample_meta("One", "Artist", "Album");
meta.path = Some("/music/One.flac".into());
let track = queries::upsert_track(&db.conn, &meta).unwrap();
let album: i64 = db
.conn
.query_row("SELECT album_id FROM tracks WHERE id = ?1", [track], |r| {
r.get(0)
})
.unwrap();
db.conn
.execute("UPDATE albums SET remote_id = 'a1' WHERE id = ?1", [album])
.unwrap();
album
}
#[test]
fn an_unstar_the_server_has_not_taken_is_not_undone() {
let dir = tempfile::tempdir().unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
let album = album(&db);
queries::queue_favourite_change(&db.conn, "album", "a1", false).unwrap();
let url = serve_with(Arc::new(Mutex::new(Vec::new())), Unstar::Unreached);
reconcile_favourites(&db, &SubsonicClient::new(&url, "u", "pw"));
let favourites = queries::favourite_album_id_set(&db.conn, queries::LOCAL_USER).unwrap();
assert!(
!favourites.contains(&album),
"the server's stale star came back"
);
assert_eq!(
queries::favourite_changes(&db.conn).unwrap().len(),
1,
"kept to send again"
);
}
#[test]
fn a_change_the_server_refuses_leaves_the_outbox() {
let dir = tempfile::tempdir().unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
queries::queue_favourite_change(&db.conn, "album", "gone", false).unwrap();
let url = serve_with(Arc::new(Mutex::new(Vec::new())), Unstar::Refused);
reconcile_favourites(&db, &SubsonicClient::new(&url, "u", "pw"));
assert!(queries::favourite_changes(&db.conn).unwrap().is_empty());
}
#[test]
fn clearing_a_sent_change_keeps_a_newer_one() {
let dir = tempfile::tempdir().unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
queries::queue_favourite_change(&db.conn, "album", "a1", true).unwrap();
let sent = queries::favourite_changes(&db.conn).unwrap().remove(0);
queries::queue_favourite_change(&db.conn, "album", "a1", false).unwrap();
queries::queue_favourite_change(&db.conn, "album", "a1", true).unwrap();
queries::forget_favourite_change(&db.conn, &sent).unwrap();
assert_eq!(queries::favourite_changes(&db.conn).unwrap().len(), 1);
}
#[test]
fn forgetting_the_server_forgets_its_changes() {
let dir = tempfile::tempdir().unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
queries::queue_favourite_change(&db.conn, "album", "a1", false).unwrap();
forget_remote(&db).unwrap();
assert!(queries::favourite_changes(&db.conn).unwrap().is_empty());
}
#[test]
fn a_change_the_server_takes_leaves_the_outbox() {
let dir = tempfile::tempdir().unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
queries::queue_favourite_change(&db.conn, "track", "s2", true).unwrap();
let stars = Arc::new(Mutex::new(Vec::new()));
let url = serve(stars.clone());
let sync = reconcile_favourites(&db, &SubsonicClient::new(&url, "u", "pw"));
assert_eq!(sync.pushed, 1);
assert_eq!(*stars.lock().unwrap(), ["s2"]);
assert!(queries::favourite_changes(&db.conn).unwrap().is_empty());
}
#[test]
fn only_the_latest_change_is_kept() {
let dir = tempfile::tempdir().unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
queries::queue_favourite_change(&db.conn, "album", "a1", false).unwrap();
queries::queue_favourite_change(&db.conn, "album", "a1", true).unwrap();
let changes = queries::favourite_changes(&db.conn).unwrap();
assert_eq!(changes.len(), 1);
assert!(changes[0].star);
}
}
#[cfg(test)]
mod refusal_tests {
use super::*;
use std::collections::HashSet;
use std::sync::Mutex;
fn serve(keys: Arc<Mutex<HashSet<String>>>) -> String {
use std::io::{BufRead, Write};
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
std::thread::spawn(move || {
for mut stream in listener.incoming().flatten() {
let mut reader = std::io::BufReader::new(stream.try_clone().unwrap());
let mut request = String::new();
reader.read_line(&mut request).unwrap();
let mut line = String::new();
while reader.read_line(&mut line).unwrap_or(0) > 2 {
line.clear();
}
let target = request.split_whitespace().nth(1).unwrap_or("");
let (path, query) = target.split_once('?').unwrap_or((target, ""));
let param = |name: &str| {
query
.split('&')
.filter_map(|kv| kv.split_once('='))
.find(|(k, _)| *k == name)
.map(|(_, v)| v.to_owned())
.unwrap_or_default()
};
let endpoint = path.rsplit('/').next().unwrap();
let body = if endpoint == "koanSignIn" {
keys.lock().unwrap().insert("second".into());
r#"{"subsonic-response":{"status":"ok","join":{"username":"mate","apiKey":"second"}}}"#.to_owned()
} else if endpoint == "getOpenSubsonicExtensions" {
r#"{"subsonic-response":{"status":"ok","openSubsonicExtensions":[{"name":"koanSignIn","versions":[1]}]}}"#.to_owned()
} else if !keys.lock().unwrap().contains(¶m("apiKey")) {
r#"{"subsonic-response":{"status":"failed","error":{"code":44,"message":"invalid API key"}}}"#.to_owned()
} else {
r#"{"subsonic-response":{"status":"ok","indexes":{"lastModified":1}}}"#
.to_owned()
};
let _ = write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nConnection: close\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
}
});
url
}
#[test]
fn a_revoked_key_is_reported_until_signing_in_again() {
let _guard = crate::config::tests::PERSIST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().unwrap();
crate::config::set_config_dir(dir.path());
let keys = Arc::new(Mutex::new(HashSet::from(["first".to_owned()])));
let url = serve(keys.clone());
Config::persist(|c| {
c.remote.enabled = true;
c.remote.url = url.clone();
c.remote.username = "mate".into();
c.remote.api_key = "first".into();
})
.unwrap();
let db = Database::open(&dir.path().join("koan.db")).unwrap();
let sync = || {
let cfg = Config::load().unwrap();
let client = SubsonicClient::from_auth(subsonic_auth(&cfg).unwrap());
let _ = sync_remote(&db, &client, Walk::IfChanged, &url, "mate", &|_| {});
};
sync();
assert_eq!(remote_problem(&Config::load().unwrap()), None);
keys.lock().unwrap().remove("first");
sync();
let cfg = Config::load().unwrap();
assert_eq!(remote_problem(&cfg).as_deref(), Some(SIGN_IN_REFUSED));
assert!(sign_in_refused(&cfg));
set_remote_credentials(&url, "mate", "hunter22").unwrap();
let cfg = Config::load().unwrap();
assert_eq!(cfg.remote.api_key, "second");
assert_eq!(remote_problem(&cfg), None, "a new credential starts clean");
sync();
assert_eq!(remote_problem(&Config::load().unwrap()), None);
}
}
#[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());
}
}