Sync listening history across devices

This commit is contained in:
Ultradesu
2026-07-27 23:15:57 +01:00
parent f55571ab34
commit 0d6c89761f
8 changed files with 205 additions and 36 deletions
Generated
+1 -1
View File
@@ -1567,7 +1567,7 @@ dependencies = [
[[package]] [[package]]
name = "furumi_tui" name = "furumi_tui"
version = "0.1.8" version = "0.1.9"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"blake3", "blake3",
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "furumi_tui" name = "furumi_tui"
version = "0.1.8" version = "0.1.9"
edition = "2024" edition = "2024"
rust-version = "1.97" rust-version = "1.97"
description = "A federated P2P player for personal music libraries" description = "A federated P2P player for personal music libraries"
+37 -15
View File
@@ -1769,6 +1769,7 @@ fn start_current_audio(
state.player.position_secs = position_secs.max(0.0); state.player.position_secs = position_secs.max(0.0);
state.player.audio_analysis = player::AudioAnalysisSnapshot::default(); state.player.audio_analysis = player::AudioAnalysisSnapshot::default();
state.player.track_started_at = Some(now_epoch_seconds()); state.player.track_started_at = Some(now_epoch_seconds());
state.player.listen_id = Some(runtime.devices.new_listen_id());
state.player.prefetched_pos = None; state.player.prefetched_pos = None;
runtime.player_start_pending = true; runtime.player_start_pending = true;
runtime.player.stop(); runtime.player.stop();
@@ -1796,16 +1797,19 @@ fn start_current_audio(
} }
// The track that was playing until now was cut short by this switch. // The track that was playing until now was cut short by this switch.
let previous_started_at = state.player.track_started_at; 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 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_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 { if state.player.playing && !same_track {
report_history( report_history(
runtime, runtime,
previous.id, &previous,
previous_listen_id.as_deref(),
state.player.track_started_at, state.player.track_started_at,
state.player.position_secs.round() as i32, (state.player.position_secs * 1_000.0).round() as i64,
false, music_dht::device_sync::ListenEndReason::Replaced,
); );
} }
same_track.then_some(previous_started_at).flatten() 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.position_secs = position_secs.max(0.0);
state.player.audio_analysis = player::AudioAnalysisSnapshot::default(); state.player.audio_analysis = player::AudioAnalysisSnapshot::default();
state.player.track_started_at = same_track_started_at.or_else(|| Some(now_epoch_seconds())); 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.player.prefetched_pos = None;
state.status_message = Some(format!("{}{}", track.title, track.artist_line())); 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. /// than 5s are noise.
fn report_history( fn report_history(
runtime: &Runtime, runtime: &Runtime,
track_id: i64, track: &crate::library::models::TrackItem,
listen_id: Option<&str>,
started_at: Option<i64>, started_at: Option<i64>,
listened: i32, listened_ms: i64,
completed: bool, ended_reason: music_dht::device_sync::ListenEndReason,
) { ) {
// Ephemeral federated tracks are not library rows; no history for them. let Some(listen_id) = listen_id else {
if listened < 5 || track_id < 0 {
return; 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 || { 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"); 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() { if let Some(finished) = state.player.current.clone() {
report_history( report_history(
runtime, runtime,
finished.id, &finished,
state.player.listen_id.as_deref(),
state.player.track_started_at, state.player.track_started_at,
finished.duration_seconds.round() as i32, (finished.duration_seconds * 1_000.0).round() as i64,
true, music_dht::device_sync::ListenEndReason::Finished,
); );
} }
if has_next { 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.current = state.player.queue.get(state.player.queue_pos).cloned();
state.player.position_secs = 0.0; state.player.position_secs = 0.0;
state.player.track_started_at = Some(now_epoch_seconds()); 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_metadata(state, runtime);
push_media_update(state, runtime, true); push_media_update(state, runtime, true);
} else { } else {
state.player.current = None; state.player.current = None;
state.player.listen_id = None;
state.player.prefetched_pos = None; state.player.prefetched_pos = None;
if let Some(effect) = update::advance_after_finish(state) { if let Some(effect) = update::advance_after_finish(state) {
perform_effect(state, runtime, effect); perform_effect(state, runtime, effect);
+3
View File
@@ -1170,6 +1170,8 @@ pub struct PlayerBar {
pub audio_analysis: crate::player::AudioAnalysisSnapshot, pub audio_analysis: crate::player::AudioAnalysisSnapshot,
/// Epoch seconds when the current track started (for history reports). /// Epoch seconds when the current track started (for history reports).
pub track_started_at: Option<i64>, pub track_started_at: Option<i64>,
/// Stable id reused for every report of the current playback session.
pub listen_id: Option<String>,
/// Queue index already enqueued in the audio thread for gapless play. /// Queue index already enqueued in the audio thread for gapless play.
pub prefetched_pos: Option<usize>, pub prefetched_pos: Option<usize>,
pub volume: u8, pub volume: u8,
@@ -1191,6 +1193,7 @@ impl Default for PlayerBar {
position_secs: 0.0, position_secs: 0.0,
audio_analysis: crate::player::AudioAnalysisSnapshot::default(), audio_analysis: crate::player::AudioAnalysisSnapshot::default(),
track_started_at: None, track_started_at: None,
listen_id: None,
prefetched_pos: None, prefetched_pos: None,
original_order: None, original_order: None,
volume: 80, volume: 80,
+1
View File
@@ -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.current = state.player.queue.get(state.player.queue_pos).cloned();
state.player.position_secs = 0.0; state.player.position_secs = 0.0;
state.player.track_started_at = None; 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); state.queue_tab.cursor = state.queue_tab.cursor.min(state.player.queue.len() - 1);
return QueueRemovalOutcome { return QueueRemovalOutcome {
restart_paused: was_loaded.then_some(was_paused), restart_paused: was_loaded.then_some(was_paused),
+60 -2
View File
@@ -12,6 +12,7 @@ use std::sync::Arc;
use std::time::Duration; use std::time::Duration;
use anyhow::{Context as _, Result}; use anyhow::{Context as _, Result};
use music_dht::device_sync::{ListenEvent, ListenTrackMetadata};
use music_dht::{ByteStream, MusicDhtService, NetworkId, PeerTicket, SecretKey, StreamAcceptor}; use music_dht::{ByteStream, MusicDhtService, NetworkId, PeerTicket, SecretKey, StreamAcceptor};
use rusqlite::{Connection, OptionalExtension, params}; use rusqlite::{Connection, OptionalExtension, params};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
@@ -21,9 +22,9 @@ use crate::app::event::AppEvent;
use crate::library::Library; use crate::library::Library;
use crate::library::models::{ArtistRef, TrackItem}; 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 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 INVITE_TTL_MS: i64 = 10 * 60 * 1000;
const PAIRING_WAIT_MS: i64 = 5 * 60 * 1000; const PAIRING_WAIT_MS: i64 = 5 * 60 * 1000;
const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1); const PAIRING_RETRY_DELAY: Duration = Duration::from_secs(1);
@@ -366,6 +367,9 @@ pub enum SyncOpPayload {
target_device_id: String, target_device_id: String,
command: PlaybackCommand, command: PlaybackCommand,
}, },
ListenRecorded {
event: ListenEvent,
},
} }
impl SyncOpPayload { impl SyncOpPayload {
@@ -637,6 +641,10 @@ impl DeviceSync {
Ok((identity.device_id, identity.name)) 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) { pub fn publish_playback(&self, mut snapshot: PlaybackSnapshot) {
if snapshot.updated_at_ms <= 0 { if snapshot.updated_at_ms <= 0 {
snapshot.updated_at_ms = now_ms(); snapshot.updated_at_ms = now_ms();
@@ -1073,6 +1081,52 @@ impl DeviceSync {
Ok(()) 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<ListenEvent> {
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<()> { pub fn record_playlist_created(&self, playlist_id: i64, title: &str) -> Result<()> {
let playlist_id = self.library.ensure_playlist_sync_id(playlist_id)?; let playlist_id = self.library.ensure_playlist_sync_id(playlist_id)?;
self.record_local_op(SyncOpPayload::PlaylistCreated { self.record_local_op(SyncOpPayload::PlaylistCreated {
@@ -1584,6 +1638,9 @@ impl DeviceSync {
self.apply_playback_command(target_device_id, command, &op.op_id)?; self.apply_playback_command(target_device_id, command, &op.op_id)?;
false false
} }
SyncOpPayload::ListenRecorded { event } => self
.library
.apply_listen_event(event, &op.origin_device_id)?,
}; };
Ok(changed) Ok(changed)
} }
@@ -3229,6 +3286,7 @@ fn payload_kind(payload: &SyncOpPayload) -> &'static str {
SyncOpPayload::DeviceTrusted { .. } => "device_trusted", SyncOpPayload::DeviceTrusted { .. } => "device_trusted",
SyncOpPayload::DeviceRevoked { .. } => "device_revoked", SyncOpPayload::DeviceRevoked { .. } => "device_revoked",
SyncOpPayload::PlaybackCommand { .. } => "playback_command", SyncOpPayload::PlaybackCommand { .. } => "playback_command",
SyncOpPayload::ListenRecorded { .. } => "listen_recorded",
} }
} }
+79 -14
View File
@@ -126,10 +126,43 @@ CREATE TABLE IF NOT EXISTS history (
completed INTEGER NOT NULL DEFAULT 0, completed INTEGER NOT NULL DEFAULT 0,
played_at TEXT NOT NULL DEFAULT (datetime('now')) 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_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_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_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_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_playlist_tracks_playlist ON playlist_tracks(playlist_id);
CREATE INDEX IF NOT EXISTS idx_fed_playlist_tracks_playlist CREATE INDEX IF NOT EXISTS idx_fed_playlist_tracks_playlist
ON fed_playlist_tracks(playlist_sync_id, position); 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.file_path, t.audio_format, t.audio_bitrate, t.audio_sample_rate,
t.audio_bit_depth, t.file_size_bytes, t.audio_bit_depth, t.file_size_bytes,
t.content_id, 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 /// Plain rows handed to the federation for publishing (see
@@ -798,9 +833,14 @@ impl Library {
|row| row.get(0), |row| row.get(0),
)?; )?;
let total_play_count: i64 = conn.query_row( let total_play_count: i64 = conn.query_row(
"SELECT COUNT(*) FROM history h "SELECT
(SELECT COUNT(*) FROM history h
WHERE h.completed = 1 AND h.track_id IN WHERE h.completed = 1 AND h.track_id IN
(SELECT track_id FROM track_artists WHERE artist_id = ?1)", (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], [id],
|row| row.get(0), |row| row.get(0),
)?; )?;
@@ -2196,20 +2236,45 @@ impl Library {
Ok(true) Ok(true)
} }
pub fn add_history( /// Idempotently materialize a portable trusted-device listen event.
pub fn apply_listen_event(
&self, &self,
track_id: i64, event: &music_dht::device_sync::ListenEvent,
started_at: Option<i64>, origin_device_id: &str,
listened_seconds: i32, ) -> Result<bool> {
completed: bool, if !event.should_record() || origin_device_id.trim().is_empty() {
) -> Result<()> { return Ok(false);
}
let content_id = music_dht::normalize_content_id(&event.content_id)
.context("invalid listen content id")?;
let conn = self.lock(); let conn = self.lock();
conn.execute( let local_track_id: Option<i64> = conn
"INSERT INTO history (track_id, started_at, listened_seconds, completed) .query_row(
VALUES (?1, ?2, ?3, ?4)", "SELECT id FROM tracks WHERE content_id = ?1 LIMIT 1",
params![track_id, started_at, listened_seconds, completed], [&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)
} }
// ----------------------------------------------------------------- // -----------------------------------------------------------------
+22 -2
View File
@@ -502,8 +502,28 @@ fn delete_track_drops_empty_release() {
fn history_counts_completed_plays() { fn history_counts_completed_plays() {
let lib = test_library(); let lib = test_library();
let track_id = add_track(&lib, "Song", "Artist", "Album"); let track_id = add_track(&lib, "Song", "Artist", "Album");
lib.add_history(track_id, None, 60, true).unwrap(); let content_id = lib
lib.add_history(track_id, None, 10, false).unwrap(); .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); let track = lib.tracks_by_ids(&[track_id]).unwrap().remove(0);
assert_eq!(track.play_count, 1); assert_eq!(track.play_count, 1);
} }