use parking_lot::Mutex;
use crate::db::connection::{Database, DbError};
use crate::db::queries::{self, HistoryCursor, OutboxEntry, OutboxKind};
use crate::remote::client::{SubsonicClient, SubsonicError};
const PAGE: u32 = 500;
const BATCH: u32 = 50;
static ONE_AT_A_TIME: Mutex<()> = Mutex::new(());
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct HistorySync {
pub sent: usize,
pub adopted: usize,
pub forgotten: usize,
pub library_synced: bool,
}
impl HistorySync {
pub fn changed(&self) -> bool {
self.adopted > 0 || self.forgotten > 0 || self.library_synced
}
}
pub fn sync(db: &Database) -> HistorySync {
let cfg = crate::config::Config::load().unwrap_or_default();
let Some(client) = crate::helpers::subsonic_client(&cfg) else {
return HistorySync::default();
};
let out = {
let _one = ONE_AT_A_TIME.lock();
let mut out = HistorySync {
sent: flush(db, &client, None),
..Default::default()
};
if let Err(e) = pull(
db,
&client,
&cfg.remote.url,
&cfg.remote.username,
true,
&mut out,
) {
log::warn!("history: could not read the server's: {e}");
}
out
};
if out.library_synced {
crate::remote::link::sync(db, crate::helpers::Walk::IfChanged);
} else if out.changed() {
crate::player::history::changed();
}
out
}
pub fn reconcile(db: &Database, client: &SubsonicClient, url: &str, username: &str) -> HistorySync {
let _one = ONE_AT_A_TIME.lock();
let mut out = HistorySync {
sent: flush(db, client, None),
..Default::default()
};
if let Err(e) = pull(db, client, url, username, false, &mut out) {
log::warn!("history: could not read the server's: {e}");
}
if out.changed() {
crate::player::history::changed();
}
out
}
pub fn scrobble(db: &Database, remote_id: &str, at_ms: i64) {
if let Err(e) = queries::queue_scrobble(&db.conn, remote_id, at_ms) {
log::warn!("history: could not queue a scrobble: {e}");
return;
}
let cfg = crate::config::Config::load().unwrap_or_default();
let Some(client) = crate::helpers::subsonic_client(&cfg) else {
return;
};
if let Some(_one) = ONE_AT_A_TIME.try_lock() {
flush(db, &client, Some(1));
}
}
pub fn forget(db: &Database, ids: &[i64]) -> Result<usize, DbError> {
if signed_in() {
for (remote_id, played_at) in
queries::remote_ids_of_plays(&db.conn, queries::LOCAL_USER, ids)?
{
queries::queue_forget(&db.conn, &remote_id, played_at * 1000)?;
}
}
let removed = queries::delete_plays(&db.conn, queries::LOCAL_USER, ids)?;
flush_in_background();
Ok(removed)
}
pub fn clear(db: &Database) -> Result<usize, DbError> {
if signed_in() {
queries::queue_clear(&db.conn, now_ms())?;
}
let removed = queries::clear_play_history(&db.conn, queries::LOCAL_USER)?;
flush_in_background();
Ok(removed)
}
fn signed_in() -> bool {
let cfg = crate::config::Config::load().unwrap_or_default();
crate::helpers::subsonic_auth(&cfg).is_some()
}
fn flush_in_background() {
let spawned = std::thread::Builder::new()
.name("koan-history-flush".into())
.spawn(|| {
let Ok(db) = crate::db::pool::shared().get() else {
return;
};
let cfg = crate::config::Config::load().unwrap_or_default();
if let Some(client) = crate::helpers::subsonic_client(&cfg) {
let _one = ONE_AT_A_TIME.lock();
flush(&db, &client, None);
}
});
if let Err(e) = spawned {
log::warn!("history: could not start sending: {e}");
}
}
fn shares_history(client: &SubsonicClient) -> Option<bool> {
crate::remote::profile::for_auth(client.auth())
.map(|p| p.offers(crate::remote::profile::HISTORY))
}
fn flush(db: &Database, client: &SubsonicClient, batches: Option<usize>) -> usize {
let mut sent = 0;
let mut requests = 0;
loop {
if batches.is_some_and(|max| requests >= max) {
break;
}
let waiting = match queries::history_outbox(&db.conn, BATCH) {
Ok(w) => w,
Err(e) => {
log::warn!("history: could not read the outbox: {e}");
break;
}
};
let Some(kind) = waiting.first().map(|e| e.kind) else {
break;
};
let run: Vec<&OutboxEntry> = waiting.iter().take_while(|e| e.kind == kind).collect();
requests += 1;
let Some(taken) = send(client, kind, &run) else {
break;
};
let ids: Vec<i64> = run.iter().map(|e| e.id).collect();
if let Err(e) = queries::drop_from_outbox(&db.conn, &ids) {
log::warn!("history: could not empty the outbox: {e}");
break;
}
sent += taken;
}
sent
}
fn send(client: &SubsonicClient, kind: OutboxKind, run: &[&OutboxEntry]) -> Option<usize> {
let pairs: Vec<(&str, i64)> = run
.iter()
.filter_map(|e| Some((e.remote_id.as_deref()?, e.at_ms)))
.collect();
match kind {
OutboxKind::Scrobble => match client.scrobble_many(&pairs) {
Ok(()) => Some(pairs.len()),
Err(SubsonicError::Api { code, .. }) if code == NOT_FOUND && pairs.len() > 1 => {
let mut taken = 0;
for pair in &pairs {
match client.scrobble_many(std::slice::from_ref(pair)) {
Ok(()) => taken += 1,
Err(e) => refused_for_good(e)?,
}
}
Some(taken)
}
Err(e) => refused_for_good(e).map(|()| 0),
},
OutboxKind::Forget | OutboxKind::Clear => {
if !shares_history(client)? {
return Some(0);
}
let result = match kind {
OutboxKind::Forget => client.koan_forget_plays(&pairs),
_ => run
.iter()
.try_for_each(|e| client.koan_forget_plays_through(e.at_ms)),
};
match result {
Ok(()) => Some(run.len()),
Err(e) => refused_for_good(e).map(|()| 0),
}
}
}
}
const NOT_FOUND: i32 = 70;
const MISSING_PARAMETER: i32 = 10;
fn refused_for_good(e: SubsonicError) -> Option<()> {
match e {
SubsonicError::Api { code, message } if matches!(code, NOT_FOUND | MISSING_PARAMETER) => {
log::info!("history: server refused for good ({code}: {message})");
Some(())
}
e => {
log::info!("history: not sent, kept for later: {e}");
None
}
}
}
fn pull(
db: &Database,
client: &SubsonicClient,
url: &str,
username: &str,
strict: bool,
out: &mut HistorySync,
) -> Result<(), SubsonicError> {
if shares_history(client) != Some(true) {
return Ok(());
}
let mut cursor = queries::history_cursor(&db.conn, url).map_err(db_failed)?;
let mut held: Option<HistoryCursor> = None;
loop {
let page = client.koan_history(cursor, PAGE)?;
let next = HistoryCursor::parse(&page.cursor).ok_or(SubsonicError::BadResponse)?;
let played: Vec<String> = page.play.iter().map(|p| p.id.clone()).collect();
let tracks = queries::track_ids_for_remote_ids(&db.conn, &played).map_err(db_failed)?;
if let Some(i) = tracks.iter().position(Option::is_none) {
if strict {
out.library_synced = true;
return Ok(());
}
held.get_or_insert(HistoryCursor {
forgotten: cursor.forgotten,
play: page.play[i].seq.map_or(cursor.play, |seq| seq - 1),
});
}
for f in &page.forgotten {
let secs = f.played / 1000;
let removed = match &f.id {
None => queries::forget_plays_through(&db.conn, queries::LOCAL_USER, secs),
Some(id) => {
let track =
queries::track_ids_for_remote_ids(&db.conn, std::slice::from_ref(id))
.map_err(db_failed)?;
match track.first().copied().flatten() {
Some(track) => {
queries::forget_play_near(&db.conn, queries::LOCAL_USER, track, secs)
}
None => Ok(0),
}
}
};
out.forgotten += removed.map_err(db_failed)?;
}
let plays: Vec<(i64, i64, Option<i64>)> = page
.play
.iter()
.zip(&tracks)
.filter_map(|(p, track)| Some(((*track)?, p.played / 1000, p.listened_ms)))
.collect();
out.adopted +=
queries::adopt_plays(&db.conn, queries::LOCAL_USER, &plays).map_err(db_failed)?;
queries::set_history_cursor(&db.conn, url, username, held.unwrap_or(next))
.map_err(db_failed)?;
cursor = next;
if !page.more {
if held.is_some() {
log::info!("history: holding the cursor before a track not synced yet");
}
return Ok(());
}
}
}
fn db_failed(e: DbError) -> SubsonicError {
SubsonicError::Io(std::io::Error::other(e.to_string()))
}
fn now_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_millis() as i64)
}
#[cfg(test)]
mod tests {
use super::*;
fn api(code: i32) -> SubsonicError {
SubsonicError::Api {
code,
message: String::new(),
}
}
#[test]
fn only_lasting_refusals_drop_the_outbox() {
assert_eq!(refused_for_good(api(NOT_FOUND)), Some(()));
assert_eq!(refused_for_good(api(MISSING_PARAMETER)), Some(()));
for code in [0, 40, 41, 42, 43, 44, 50] {
assert_eq!(refused_for_good(api(code)), None, "code {code}");
}
assert_eq!(refused_for_good(SubsonicError::BadResponse), None);
}
}