diff --git a/Cargo.lock b/Cargo.lock index 5a08396..86f25ee 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1567,7 +1567,7 @@ dependencies = [ [[package]] name = "furumi_tui" -version = "0.1.8" +version = "0.1.9" dependencies = [ "anyhow", "blake3", diff --git a/Cargo.toml b/Cargo.toml index bd0d8d0..2a344f8 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "furumi_tui" -version = "0.1.8" +version = "0.1.9" edition = "2024" rust-version = "1.97" description = "A federated P2P player for personal music libraries" diff --git a/src/app/mod.rs b/src/app/mod.rs index 3990192..5d5e2f5 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -1769,6 +1769,7 @@ fn start_current_audio( state.player.position_secs = position_secs.max(0.0); state.player.audio_analysis = player::AudioAnalysisSnapshot::default(); state.player.track_started_at = Some(now_epoch_seconds()); + state.player.listen_id = Some(runtime.devices.new_listen_id()); state.player.prefetched_pos = None; runtime.player_start_pending = true; runtime.player.stop(); @@ -1796,16 +1797,19 @@ fn start_current_audio( } // The track that was playing until now was cut short by this switch. let previous_started_at = state.player.track_started_at; + let previous_listen_id = state.player.listen_id.take(); let next_key = track_playback_key(&track); + let mut same_track = false; let same_track_started_at = if let Some(previous) = state.player.current.take() { - let same_track = track_playback_key(&previous) == next_key; + same_track = track_playback_key(&previous) == next_key; if state.player.playing && !same_track { report_history( runtime, - previous.id, + &previous, + previous_listen_id.as_deref(), state.player.track_started_at, - state.player.position_secs.round() as i32, - false, + (state.player.position_secs * 1_000.0).round() as i64, + music_dht::device_sync::ListenEndReason::Replaced, ); } same_track.then_some(previous_started_at).flatten() @@ -1818,6 +1822,11 @@ fn start_current_audio( state.player.position_secs = position_secs.max(0.0); state.player.audio_analysis = player::AudioAnalysisSnapshot::default(); state.player.track_started_at = same_track_started_at.or_else(|| Some(now_epoch_seconds())); + state.player.listen_id = if same_track { + previous_listen_id + } else { + Some(runtime.devices.new_listen_id()) + }; state.player.prefetched_pos = None; state.status_message = Some(format!("▶ {} — {}", track.title, track.artist_line())); @@ -1972,18 +1981,28 @@ fn maybe_prefetch_next(state: &mut AppState, runtime: &Runtime) { /// than 5s are noise. fn report_history( runtime: &Runtime, - track_id: i64, + track: &crate::library::models::TrackItem, + listen_id: Option<&str>, started_at: Option, - listened: i32, - completed: bool, + listened_ms: i64, + ended_reason: music_dht::device_sync::ListenEndReason, ) { - // Ephemeral federated tracks are not library rows; no history for them. - if listened < 5 || track_id < 0 { + let Some(listen_id) = listen_id else { return; - } - let library = Arc::clone(&runtime.library); + }; + let Some(event) = runtime.devices.listen_event_for_track( + listen_id.to_string(), + track, + started_at.unwrap_or_else(now_epoch_seconds) * 1_000, + listened_ms, + ended_reason, + ) else { + tracing::warn!(title = %track.title, "history skipped: track has no content id"); + return; + }; + let devices = Arc::clone(&runtime.devices); tokio::task::spawn_blocking(move || { - if let Err(err) = library.add_history(track_id, started_at, listened, completed) { + if let Err(err) = devices.record_listen(event) { tracing::warn!(%err, "history write failed"); } }); @@ -3292,10 +3311,11 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent if let Some(finished) = state.player.current.clone() { report_history( runtime, - finished.id, + &finished, + state.player.listen_id.as_deref(), state.player.track_started_at, - finished.duration_seconds.round() as i32, - true, + (finished.duration_seconds * 1_000.0).round() as i64, + music_dht::device_sync::ListenEndReason::Finished, ); } if has_next { @@ -3309,10 +3329,12 @@ fn handle_app_event(state: &mut AppState, runtime: &mut Runtime, event: AppEvent state.player.current = state.player.queue.get(state.player.queue_pos).cloned(); state.player.position_secs = 0.0; state.player.track_started_at = Some(now_epoch_seconds()); + state.player.listen_id = Some(runtime.devices.new_listen_id()); push_media_metadata(state, runtime); push_media_update(state, runtime, true); } else { state.player.current = None; + state.player.listen_id = None; state.player.prefetched_pos = None; if let Some(effect) = update::advance_after_finish(state) { perform_effect(state, runtime, effect); diff --git a/src/app/state.rs b/src/app/state.rs index 06567ee..d428737 100644 --- a/src/app/state.rs +++ b/src/app/state.rs @@ -1170,6 +1170,8 @@ pub struct PlayerBar { pub audio_analysis: crate::player::AudioAnalysisSnapshot, /// Epoch seconds when the current track started (for history reports). pub track_started_at: Option, + /// Stable id reused for every report of the current playback session. + pub listen_id: Option, /// Queue index already enqueued in the audio thread for gapless play. pub prefetched_pos: Option, pub volume: u8, @@ -1191,6 +1193,7 @@ impl Default for PlayerBar { position_secs: 0.0, audio_analysis: crate::player::AudioAnalysisSnapshot::default(), track_started_at: None, + listen_id: None, prefetched_pos: None, original_order: None, volume: 80, diff --git a/src/app/update.rs b/src/app/update.rs index ee39cb8..b44c14b 100644 --- a/src/app/update.rs +++ b/src/app/update.rs @@ -1092,6 +1092,7 @@ fn remove_queue_indices(state: &mut AppState, indices: &[usize]) -> QueueRemoval state.player.current = state.player.queue.get(state.player.queue_pos).cloned(); state.player.position_secs = 0.0; state.player.track_started_at = None; + state.player.listen_id = None; state.queue_tab.cursor = state.queue_tab.cursor.min(state.player.queue.len() - 1); return QueueRemovalOutcome { restart_paused: was_loaded.then_some(was_paused), diff --git a/src/devices.rs b/src/devices.rs index 7fca063..626d436 100644 --- a/src/devices.rs +++ b/src/devices.rs @@ -12,6 +12,7 @@ use std::sync::Arc; use std::time::Duration; use anyhow::{Context as _, Result}; +use music_dht::device_sync::{ListenEvent, ListenTrackMetadata}; use music_dht::{ByteStream, MusicDhtService, NetworkId, PeerTicket, SecretKey, StreamAcceptor}; use rusqlite::{Connection, OptionalExtension, params}; use serde::{Deserialize, Serialize}; @@ -21,9 +22,9 @@ use crate::app::event::AppEvent; use crate::library::Library; use crate::library::models::{ArtistRef, TrackItem}; -pub const SYNC_ALPN: &[u8] = b"furumi/sync/1"; +pub const SYNC_ALPN: &[u8] = b"furumi/sync/2"; const CLIENT_VERSION: &str = env!("CARGO_PKG_VERSION"); -const PROTOCOL_VERSION: u16 = 1; +const PROTOCOL_VERSION: u16 = 2; const INVITE_TTL_MS: i64 = 10 * 60 * 1000; const PAIRING_WAIT_MS: i64 = 5 * 60 * 1000; const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1); @@ -366,6 +367,9 @@ pub enum SyncOpPayload { target_device_id: String, command: PlaybackCommand, }, + ListenRecorded { + event: ListenEvent, + }, } impl SyncOpPayload { @@ -637,6 +641,10 @@ impl DeviceSync { Ok((identity.device_id, identity.name)) } + pub fn new_listen_id(&self) -> String { + format!("{}-{}", now_ms(), random_hex(12)) + } + pub fn publish_playback(&self, mut snapshot: PlaybackSnapshot) { if snapshot.updated_at_ms <= 0 { snapshot.updated_at_ms = now_ms(); @@ -1073,6 +1081,52 @@ impl DeviceSync { Ok(()) } + pub fn record_listen(&self, event: ListenEvent) -> Result<()> { + if !event.should_record() { + return Ok(()); + } + self.record_local_op(SyncOpPayload::ListenRecorded { event }) + } + + pub fn listen_event_for_track( + &self, + listen_id: String, + track: &TrackItem, + started_at_ms: i64, + listened_ms: i64, + ended_reason: music_dht::device_sync::ListenEndReason, + ) -> Option { + let content_id = track + .content_id + .as_deref() + .or_else(|| track.fed.as_ref()?.content_id.as_deref()) + .and_then(music_dht::normalize_content_id)?; + Some(ListenEvent { + listen_id, + content_id, + started_at_ms, + listened_ms, + track_duration_ms: (track.duration_seconds > 0.0) + .then_some((track.duration_seconds * 1_000.0).round() as i64), + ended_reason, + track: ListenTrackMetadata { + title: track.title.clone(), + artist_names: track + .artists + .iter() + .map(|artist| artist.name.clone()) + .collect(), + featured_artist_names: track + .featured_artists + .iter() + .map(|artist| artist.name.clone()) + .collect(), + release_title: (!track.release_title.trim().is_empty()) + .then(|| track.release_title.clone()), + }, + }) + } + pub fn record_playlist_created(&self, playlist_id: i64, title: &str) -> Result<()> { let playlist_id = self.library.ensure_playlist_sync_id(playlist_id)?; self.record_local_op(SyncOpPayload::PlaylistCreated { @@ -1584,6 +1638,9 @@ impl DeviceSync { self.apply_playback_command(target_device_id, command, &op.op_id)?; false } + SyncOpPayload::ListenRecorded { event } => self + .library + .apply_listen_event(event, &op.origin_device_id)?, }; Ok(changed) } @@ -3229,6 +3286,7 @@ fn payload_kind(payload: &SyncOpPayload) -> &'static str { SyncOpPayload::DeviceTrusted { .. } => "device_trusted", SyncOpPayload::DeviceRevoked { .. } => "device_revoked", SyncOpPayload::PlaybackCommand { .. } => "playback_command", + SyncOpPayload::ListenRecorded { .. } => "listen_recorded", } } diff --git a/src/library/mod.rs b/src/library/mod.rs index d31137b..e913cf3 100644 --- a/src/library/mod.rs +++ b/src/library/mod.rs @@ -126,10 +126,43 @@ CREATE TABLE IF NOT EXISTS history ( completed INTEGER NOT NULL DEFAULT 0, played_at TEXT NOT NULL DEFAULT (datetime('now')) ); +CREATE TABLE IF NOT EXISTS listen_events ( + listen_id TEXT PRIMARY KEY, + content_id TEXT NOT NULL, + local_track_id INTEGER REFERENCES tracks(id) ON DELETE SET NULL, + origin_device_id TEXT NOT NULL, + started_at_ms INTEGER NOT NULL, + listened_ms INTEGER NOT NULL, + track_duration_ms INTEGER, + ended_reason TEXT NOT NULL, + qualified INTEGER NOT NULL, + metadata_json TEXT NOT NULL, + created_at TEXT NOT NULL DEFAULT (datetime('now')) +); CREATE INDEX IF NOT EXISTS idx_tracks_release ON tracks(release_id); CREATE INDEX IF NOT EXISTS idx_track_artists_artist ON track_artists(artist_id); CREATE INDEX IF NOT EXISTS idx_release_artists_artist ON release_artists(artist_id); CREATE INDEX IF NOT EXISTS idx_history_track ON history(track_id); +CREATE INDEX IF NOT EXISTS idx_listen_events_content + ON listen_events(content_id, started_at_ms DESC); +CREATE INDEX IF NOT EXISTS idx_listen_events_local_track + ON listen_events(local_track_id, qualified); +CREATE TRIGGER IF NOT EXISTS reconcile_listen_events_after_track_insert +AFTER INSERT ON tracks +WHEN NEW.content_id IS NOT NULL +BEGIN + UPDATE listen_events + SET local_track_id = NEW.id + WHERE content_id = NEW.content_id; +END; +CREATE TRIGGER IF NOT EXISTS reconcile_listen_events_after_track_content_id +AFTER UPDATE OF content_id ON tracks +WHEN NEW.content_id IS NOT NULL +BEGIN + UPDATE listen_events + SET local_track_id = NEW.id + WHERE content_id = NEW.content_id; +END; CREATE INDEX IF NOT EXISTS idx_playlist_tracks_playlist ON playlist_tracks(playlist_id); CREATE INDEX IF NOT EXISTS idx_fed_playlist_tracks_playlist ON fed_playlist_tracks(playlist_sync_id, position); @@ -159,7 +192,9 @@ const TRACK_COLUMNS: &str = " t.file_path, t.audio_format, t.audio_bitrate, t.audio_sample_rate, t.audio_bit_depth, t.file_size_bytes, t.content_id, - (SELECT COUNT(*) FROM history h WHERE h.track_id = t.id AND h.completed = 1) + ((SELECT COUNT(*) FROM history h WHERE h.track_id = t.id AND h.completed = 1) + + (SELECT COUNT(*) FROM listen_events le + WHERE le.local_track_id = t.id AND le.qualified = 1)) "; /// Plain rows handed to the federation for publishing (see @@ -798,9 +833,14 @@ impl Library { |row| row.get(0), )?; let total_play_count: i64 = conn.query_row( - "SELECT COUNT(*) FROM history h - WHERE h.completed = 1 AND h.track_id IN - (SELECT track_id FROM track_artists WHERE artist_id = ?1)", + "SELECT + (SELECT COUNT(*) FROM history h + WHERE h.completed = 1 AND h.track_id IN + (SELECT track_id FROM track_artists WHERE artist_id = ?1)) + + + (SELECT COUNT(*) FROM listen_events le + WHERE le.qualified = 1 AND le.local_track_id IN + (SELECT track_id FROM track_artists WHERE artist_id = ?1))", [id], |row| row.get(0), )?; @@ -2196,20 +2236,45 @@ impl Library { Ok(true) } - pub fn add_history( + /// Idempotently materialize a portable trusted-device listen event. + pub fn apply_listen_event( &self, - track_id: i64, - started_at: Option, - listened_seconds: i32, - completed: bool, - ) -> Result<()> { + event: &music_dht::device_sync::ListenEvent, + origin_device_id: &str, + ) -> Result { + if !event.should_record() || origin_device_id.trim().is_empty() { + return Ok(false); + } + let content_id = music_dht::normalize_content_id(&event.content_id) + .context("invalid listen content id")?; let conn = self.lock(); - conn.execute( - "INSERT INTO history (track_id, started_at, listened_seconds, completed) - VALUES (?1, ?2, ?3, ?4)", - params![track_id, started_at, listened_seconds, completed], + let local_track_id: Option = conn + .query_row( + "SELECT id FROM tracks WHERE content_id = ?1 LIMIT 1", + [&content_id], + |row| row.get(0), + ) + .optional()?; + let inserted = conn.execute( + "INSERT OR IGNORE INTO listen_events + (listen_id, content_id, local_track_id, origin_device_id, + started_at_ms, listened_ms, track_duration_ms, ended_reason, + qualified, metadata_json) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", + params![ + event.listen_id, + content_id, + local_track_id, + origin_device_id, + event.started_at_ms, + event.listened_ms, + event.track_duration_ms, + serde_json::to_string(&event.ended_reason)?, + i64::from(event.qualifies_as_play()), + serde_json::to_string(&event.track)?, + ], )?; - Ok(()) + Ok(inserted > 0) } // ----------------------------------------------------------------- diff --git a/src/library/tests.rs b/src/library/tests.rs index 958ca18..57e1d44 100644 --- a/src/library/tests.rs +++ b/src/library/tests.rs @@ -502,8 +502,28 @@ fn delete_track_drops_empty_release() { fn history_counts_completed_plays() { let lib = test_library(); let track_id = add_track(&lib, "Song", "Artist", "Album"); - lib.add_history(track_id, None, 60, true).unwrap(); - lib.add_history(track_id, None, 10, false).unwrap(); + let content_id = lib + .tracks_by_ids(&[track_id]) + .unwrap() + .remove(0) + .content_id + .unwrap(); + let event = music_dht::device_sync::ListenEvent { + listen_id: "listen-1".to_string(), + content_id, + started_at_ms: 1_700_000_000_000, + listened_ms: 60_000, + track_duration_ms: Some(60_000), + ended_reason: music_dht::device_sync::ListenEndReason::Finished, + track: music_dht::device_sync::ListenTrackMetadata { + title: "Song".to_string(), + artist_names: vec!["Artist".to_string()], + featured_artist_names: Vec::new(), + release_title: Some("Album".to_string()), + }, + }; + assert!(lib.apply_listen_event(&event, "device-a").unwrap()); + assert!(!lib.apply_listen_event(&event, "device-a").unwrap()); let track = lib.tracks_by_ids(&[track_id]).unwrap().remove(0); assert_eq!(track.play_count, 1); }