use std::thread;
use std::time::{Duration, Instant};
use crossbeam_channel::{RecvTimeoutError, Sender, TrySendError};
use crate::db::connection::Database;
use crate::db::queries;
use crate::player::state::QueueItemId;
use crate::remote::client::PlaybackReportState;
const MAX_TICK_MS: u64 = 2_000;
const SCROBBLE_MIN_TRACK_MS: u64 = 30_000;
const SCROBBLE_ENOUGH_MS: u64 = 4 * 60_000;
const QUEUE_DEPTH: usize = 64;
const REPORT_INTERVAL: Duration = Duration::from_secs(1);
#[derive(Debug, Clone)]
pub struct InFlight {
pub item: QueueItemId,
track_id: Option<i64>,
last_position_ms: u64,
listened_ms: u64,
}
impl InFlight {
pub fn new(item: QueueItemId, track_id: Option<i64>) -> Self {
Self {
item,
track_id,
last_position_ms: 0,
listened_ms: 0,
}
}
pub fn track_id(&self) -> Option<i64> {
self.track_id
}
#[cfg(test)]
pub fn track_id_for_test(&mut self, track_id: i64) {
self.track_id = Some(track_id);
}
pub fn advance(&mut self, position_ms: u64) {
let delta = position_ms.saturating_sub(self.last_position_ms);
if delta <= MAX_TICK_MS {
self.listened_ms += delta;
}
self.last_position_ms = position_ms;
}
pub fn listened_ms(&self) -> u64 {
self.listened_ms
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PlaybackReport {
pub track_id: i64,
pub state: PlaybackReportState,
pub position_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PlayEvent {
Started { track_id: i64, position_ms: u64 },
Playback(PlaybackReport),
Finished { track_id: i64, listened_ms: u64 },
}
pub struct PlayRecorder {
tx: Sender<PlayEvent>,
}
impl PlayRecorder {
pub fn spawn() -> Option<Self> {
let db = match Database::open_default() {
Ok(db) => db,
Err(e) => {
log::warn!("play history disabled — cannot open database: {e}");
return None;
}
};
let (tx, rx) = crossbeam_channel::bounded::<PlayEvent>(QUEUE_DEPTH);
thread::Builder::new()
.name("koan-history".into())
.spawn(move || {
let mut writer = Writer::new(db);
loop {
let next = match writer.reports.deadline() {
Some(deadline) => rx.recv_deadline(deadline),
None => rx.recv().map_err(|_| RecvTimeoutError::Disconnected),
};
match next {
Ok(event) => writer.handle(event),
Err(RecvTimeoutError::Timeout) => writer.send_due_report(),
Err(RecvTimeoutError::Disconnected) => break,
}
}
})
.map_err(|e| log::warn!("play history disabled — cannot spawn writer: {e}"))
.ok()?;
Some(Self { tx })
}
#[cfg(test)]
pub fn capture() -> (Self, crossbeam_channel::Receiver<PlayEvent>) {
let (tx, rx) = crossbeam_channel::bounded(QUEUE_DEPTH);
(Self { tx }, rx)
}
pub fn record(&self, event: PlayEvent) {
match self.tx.try_send(event) {
Ok(()) => {}
Err(TrySendError::Full(_)) => {
log::warn!("play history writer is behind — dropping {event:?}")
}
Err(TrySendError::Disconnected(_)) => {}
}
}
}
struct Writer {
db: Database,
open: Option<(i64, i64)>,
started: Option<(i64, u64)>,
reports: Coalescer,
}
impl Writer {
fn new(db: Database) -> Self {
Self {
db,
open: None,
started: None,
reports: Coalescer::default(),
}
}
fn handle(&mut self, event: PlayEvent) {
match event {
PlayEvent::Started {
track_id,
position_ms,
} => {
match queries::record_play(&self.db.conn, queries::LOCAL_USER, track_id, None) {
Ok(id) => self.open = Some((id, track_id)),
Err(e) => {
self.open = None;
log::warn!("failed to record play of track {track_id}: {e}");
}
}
self.started = Some((track_id, now_ms()));
self.report(PlaybackReport {
track_id,
state: PlaybackReportState::Playing,
position_ms,
});
}
PlayEvent::Playback(report) => self.report(report),
PlayEvent::Finished {
track_id,
listened_ms,
} => {
if let Some((started_track, at_ms)) = self.started.take()
&& started_track == track_id
{
scrobble_if_heard(&self.db, track_id, listened_ms, at_ms);
}
let Some((id, started)) = self.open.take() else {
return;
};
if started != track_id {
return;
}
if let Err(e) =
queries::set_listened_ms(&self.db.conn, id, track_id, listened_ms as i64)
{
log::warn!("failed to record listening time for track {track_id}: {e}");
}
}
}
}
fn report(&mut self, report: PlaybackReport) {
if let Some(report) = self.reports.offer(report, Instant::now()) {
send_report(&self.db, report);
}
}
fn send_due_report(&mut self) {
if let Some(report) = self.reports.due(Instant::now()) {
send_report(&self.db, report);
}
}
}
#[derive(Debug, Default)]
struct Coalescer {
last_sent: Option<Instant>,
pending: Option<PlaybackReport>,
}
impl Coalescer {
fn offer(&mut self, report: PlaybackReport, now: Instant) -> Option<PlaybackReport> {
let recent = self
.last_sent
.is_some_and(|at| now.duration_since(at) < REPORT_INTERVAL);
if report.state == PlaybackReportState::Playing && recent {
self.pending = Some(report);
return None;
}
self.pending = None;
self.last_sent = Some(now);
Some(report)
}
fn deadline(&self) -> Option<Instant> {
self.pending?;
Some(self.last_sent? + REPORT_INTERVAL)
}
fn due(&mut self, now: Instant) -> Option<PlaybackReport> {
if self.deadline()? > now {
return None;
}
self.last_sent = Some(now);
self.pending.take()
}
}
fn now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_millis() as u64)
}
fn counts_as_heard(listened_ms: u64, duration_ms: Option<u64>) -> bool {
match duration_ms.filter(|&d| d > 0) {
Some(d) => d >= SCROBBLE_MIN_TRACK_MS && listened_ms >= (d / 2).min(SCROBBLE_ENOUGH_MS),
None => listened_ms >= SCROBBLE_ENOUGH_MS,
}
}
fn remote_track(db: &Database, track_id: i64) -> Option<(String, Option<u64>)> {
let track = queries::get_track_row(&db.conn, track_id).ok()??;
let duration_ms = track.duration_ms.map(|d| d.max(0) as u64);
Some((track.remote_id?, duration_ms))
}
fn remote_client() -> Option<std::sync::Arc<crate::remote::client::SubsonicClient>> {
let cfg = crate::config::Config::load().unwrap_or_default();
crate::helpers::subsonic_client(&cfg)
}
fn send_report(db: &Database, report: PlaybackReport) {
let Some((remote_id, _)) = remote_track(db, report.track_id) else {
return;
};
let Some(client) = remote_client() else {
return;
};
if let Err(e) = client.report_playback(&remote_id, report.state, report.position_ms) {
log::warn!(
"failed to report playback of track {} to remote: {e}",
report.track_id
);
}
}
fn scrobble_if_heard(db: &Database, track_id: i64, listened_ms: u64, at_ms: u64) {
let Some((remote_id, duration_ms)) = remote_track(db, track_id) else {
return;
};
if !counts_as_heard(listened_ms, duration_ms) {
return;
}
let Some(client) = remote_client() else {
return;
};
if let Err(e) = client.scrobble(&remote_id, at_ms) {
log::warn!("failed to report track {track_id} to remote: {e}");
}
}
#[cfg(test)]
mod tests {
use super::*;
fn report(state: PlaybackReportState, position_ms: u64) -> PlaybackReport {
PlaybackReport {
track_id: 1,
state,
position_ms,
}
}
#[test]
fn a_seek_bar_drag_sends_where_it_lands() {
use PlaybackReportState::Playing;
let mut c = Coalescer::default();
let t0 = Instant::now();
assert_eq!(c.offer(report(Playing, 0), t0), Some(report(Playing, 0)));
for step in 1..=20 {
let at = t0 + Duration::from_millis(step * 30);
assert_eq!(c.offer(report(Playing, step * 5_000), at), None);
}
assert_eq!(c.deadline(), Some(t0 + REPORT_INTERVAL));
assert_eq!(c.due(t0 + Duration::from_millis(900)), None, "not yet");
assert_eq!(
c.due(t0 + REPORT_INTERVAL),
Some(report(Playing, 100_000)),
"only the last step"
);
assert_eq!(c.deadline(), None);
}
#[test]
fn playing_after_a_quiet_second_goes_straight_out() {
use PlaybackReportState::Playing;
let mut c = Coalescer::default();
let t0 = Instant::now();
c.offer(report(Playing, 0), t0);
assert_eq!(
c.offer(report(Playing, 60_000), t0 + REPORT_INTERVAL),
Some(report(Playing, 60_000))
);
}
#[test]
fn a_pause_or_stop_is_never_held_and_supersedes_a_held_seek() {
use PlaybackReportState::{Paused, Playing, Stopped};
let mut c = Coalescer::default();
let t0 = Instant::now();
c.offer(report(Playing, 0), t0);
assert_eq!(c.offer(report(Playing, 30_000), t0), None);
assert_eq!(
c.offer(report(Paused, 30_000), t0),
Some(report(Paused, 30_000))
);
assert_eq!(c.deadline(), None, "the held seek is gone");
assert_eq!(
c.offer(report(Stopped, 30_000), t0),
Some(report(Stopped, 30_000))
);
}
fn flight() -> InFlight {
InFlight::new(QueueItemId::new(), Some(1))
}
#[test]
fn listening_accumulates_across_ticks() {
let mut f = flight();
for tick in 1..=100 {
f.advance(tick * 50);
}
assert_eq!(f.listened_ms(), 5_000);
}
#[test]
fn a_pause_contributes_nothing() {
let mut f = flight();
f.advance(1_000);
for _ in 0..100 {
f.advance(1_000);
}
assert_eq!(f.listened_ms(), 1_000);
}
#[test]
fn seeking_forward_does_not_credit_the_skipped_stretch() {
let mut f = flight();
f.advance(1_000);
f.advance(280_000); f.advance(280_050);
assert_eq!(f.listened_ms(), 1_050);
}
#[test]
fn seeking_backward_does_not_go_negative_or_double_count() {
let mut f = flight();
f.advance(100_000);
f.advance(0); f.advance(50);
assert_eq!(f.listened_ms(), 50);
}
#[test]
fn heard_means_half_the_track() {
assert!(!counts_as_heard(89_999, Some(180_000)));
assert!(counts_as_heard(90_000, Some(180_000)));
}
#[test]
fn four_minutes_is_enough_for_a_long_track() {
assert!(counts_as_heard(240_000, Some(20 * 60_000)));
assert!(!counts_as_heard(239_999, Some(20 * 60_000)));
}
#[test]
fn a_track_under_thirty_seconds_never_counts() {
assert!(!counts_as_heard(29_000, Some(29_000)));
}
#[test]
fn a_skip_does_not_count() {
assert!(!counts_as_heard(2_000, Some(200_000)));
}
#[test]
fn unknown_duration_needs_four_minutes() {
assert!(!counts_as_heard(120_000, None));
assert!(counts_as_heard(240_000, None));
}
#[test]
fn starting_mid_track_does_not_credit_the_offset() {
let mut f = flight();
f.advance(120_000);
f.advance(120_050);
assert_eq!(f.listened_ms(), 50);
}
}