From 682794755ee083e07555bbc8e98006fc1e05d1c1 Mon Sep 17 00:00:00 2001 From: Ultradesu Date: Fri, 24 Jul 2026 03:06:20 +0300 Subject: [PATCH] Connected Devices: fixed rplaylist sync --- src/app/mod.rs | 9 +- src/app/update.rs | 11 ++ src/devices.rs | 247 +++++++++++++++++++++++++++++++++++++++++---- src/library/mod.rs | 64 ++++++++---- 4 files changed, 291 insertions(+), 40 deletions(-) diff --git a/src/app/mod.rs b/src/app/mod.rs index 9bebf91..109d5f3 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -481,15 +481,20 @@ fn perform_effect(state: &mut AppState, runtime: &mut Runtime, effect: Effect) { Effect::RemoveFromPlaylist { playlist_id, track_ids, + content_ids, } => { let library = Arc::clone(&runtime.library); let devices = Arc::clone(&runtime.devices); let tx = runtime.event_tx.clone(); tokio::task::spawn_blocking(move || { - let event = match library.remove_tracks_from_playlist(playlist_id, &track_ids) { + let event = match library + .remove_tracks_from_playlist(playlist_id, &track_ids) + .and_then(|()| { + library.remove_content_ids_from_playlist(playlist_id, &content_ids) + }) { Ok(()) => { if let Err(err) = - devices.record_playlist_tracks_removed(playlist_id, &track_ids) + devices.record_playlist_content_removed(playlist_id, &content_ids) { tracing::warn!(%err, playlist_id, "recording synced playlist removal failed"); } diff --git a/src/app/update.rs b/src/app/update.rs index fa8dc8c..039eaff 100644 --- a/src/app/update.rs +++ b/src/app/update.rs @@ -37,6 +37,7 @@ pub enum Effect { RemoveFromPlaylist { playlist_id: i64, track_ids: Vec, + content_ids: Vec, }, RemoveQueueIndices { indices: Vec, @@ -607,6 +608,15 @@ fn delete_selected(state: &mut AppState) -> Option { return None; } let track_ids: Vec = tracks.iter().map(|track| track.id).collect(); + let content_ids: Vec = tracks + .iter() + .filter_map(|track| { + track + .content_id + .as_deref() + .and_then(music_dht::normalize_content_id) + }) + .collect(); state.track_selection.clear(); if opened.id == super::state::LIKES_PLAYLIST_ID { let liked: Vec = track_ids @@ -626,6 +636,7 @@ fn delete_selected(state: &mut AppState) -> Option { return Some(Effect::RemoveFromPlaylist { playlist_id: opened.id, track_ids, + content_ids, }); } } diff --git a/src/devices.rs b/src/devices.rs index c6802cb..b79a32f 100644 --- a/src/devices.rs +++ b/src/devices.rs @@ -5,7 +5,7 @@ //! materialized tables, so offline clients can merge likes, playlists and //! membership changes deterministically. -use std::collections::BTreeMap; +use std::collections::{BTreeMap, BTreeSet}; use std::path::PathBuf; use std::str::FromStr; use std::sync::Arc; @@ -191,7 +191,13 @@ struct SyncSnapshot { #[serde(default)] likes: Vec, #[serde(default)] + unlikes: Vec, + #[serde(default)] playlists: Vec, + #[serde(default)] + deleted_playlists: Vec, + #[serde(default)] + removed_playlist_items: Vec, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -203,6 +209,13 @@ struct SnapshotLike { fed: Option, } +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SnapshotLikeTombstone { + content_id: String, + hlc_ms: i64, + op_id: String, +} + #[derive(Debug, Clone, Serialize, Deserialize)] struct SnapshotPlaylist { playlist_id: String, @@ -213,6 +226,13 @@ struct SnapshotPlaylist { items: Vec, } +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SnapshotPlaylistTombstone { + playlist_id: String, + hlc_ms: i64, + op_id: String, +} + #[derive(Debug, Clone, Serialize, Deserialize)] struct SnapshotPlaylistItem { content_id: String, @@ -223,6 +243,14 @@ struct SnapshotPlaylistItem { fed: Option, } +#[derive(Debug, Clone, Serialize, Deserialize)] +struct SnapshotPlaylistItemTombstone { + playlist_id: String, + content_id: String, + hlc_ms: i64, + op_id: String, +} + enum PairAttempt { Accepted(String), Pending, @@ -769,15 +797,22 @@ impl DeviceSync { Ok(()) } - pub fn record_playlist_tracks_removed( + pub fn record_playlist_content_removed( &self, playlist_id: i64, - track_ids: &[i64], + content_ids: &[String], ) -> Result<()> { let Some(playlist_id) = self.library.playlist_sync_id(playlist_id)? else { return Ok(()); }; - for content_id in self.library.track_content_ids(track_ids)? { + let mut seen = BTreeSet::new(); + for content_id in content_ids { + let Some(content_id) = music_dht::normalize_content_id(content_id) else { + continue; + }; + if !seen.insert(content_id.clone()) { + continue; + } self.record_local_op(SyncOpPayload::PlaylistTrackRemoved { playlist_id: playlist_id.clone(), content_id, @@ -1442,6 +1477,10 @@ impl DeviceSync { &like.op_id, )?; } + for like in snapshot.unlikes { + changed |= + self.apply_like_state(&like.content_id, false, None, like.hlc_ms, &like.op_id)?; + } for playlist in snapshot.playlists { changed |= self.apply_playlist_state( &playlist.playlist_id, @@ -1462,6 +1501,26 @@ impl DeviceSync { )?; } } + for playlist in snapshot.deleted_playlists { + changed |= self.apply_playlist_state( + &playlist.playlist_id, + "", + true, + playlist.hlc_ms, + &playlist.op_id, + )?; + } + for item in snapshot.removed_playlist_items { + changed |= self.apply_playlist_item_state( + &item.playlist_id, + &item.content_id, + false, + 0, + None, + item.hlc_ms, + &item.op_id, + )?; + } if changed { self.notify_library_changed(); } @@ -1636,6 +1695,23 @@ impl DeviceSync { fed, }); } + let unlikes = { + let conn = lock(&self.conn); + let mut stmt = conn.prepare( + "SELECT content_id, hlc_ms, op_id + FROM sync_state_likes + WHERE liked = 0 + ORDER BY content_id", + )?; + stmt.query_map([], |row| { + Ok(SnapshotLikeTombstone { + content_id: row.get(0)?, + hlc_ms: row.get(1)?, + op_id: row.get(2)?, + }) + })? + .collect::>>()? + }; let conn = lock(&self.conn); let mut playlist_stmt = conn.prepare( @@ -1698,7 +1774,46 @@ impl DeviceSync { items, }); } - Ok(SyncSnapshot { likes, playlists }) + let deleted_playlists = { + let mut stmt = conn.prepare( + "SELECT playlist_id, hlc_ms, op_id + FROM sync_state_playlists + WHERE deleted = 1 + ORDER BY playlist_id", + )?; + stmt.query_map([], |row| { + Ok(SnapshotPlaylistTombstone { + playlist_id: row.get(0)?, + hlc_ms: row.get(1)?, + op_id: row.get(2)?, + }) + })? + .collect::>>()? + }; + let removed_playlist_items = { + let mut stmt = conn.prepare( + "SELECT playlist_id, content_id, hlc_ms, op_id + FROM sync_state_playlist_items + WHERE present = 0 + ORDER BY playlist_id, content_id", + )?; + stmt.query_map([], |row| { + Ok(SnapshotPlaylistItemTombstone { + playlist_id: row.get(0)?, + content_id: row.get(1)?, + hlc_ms: row.get(2)?, + op_id: row.get(3)?, + }) + })? + .collect::>>()? + }; + Ok(SyncSnapshot { + likes, + unlikes, + playlists, + deleted_playlists, + removed_playlist_items, + }) } fn device_profiles(&self) -> Result> { @@ -2548,6 +2663,7 @@ fn payload_kind(payload: &SyncOpPayload) -> &'static str { } fn compactable_tombstone_ids(conn: &Connection) -> Result> { + let own_device_id = get_meta(conn, "device_id")?.unwrap_or_default(); let active_devices = conn .prepare( "SELECT device_id FROM sync_devices @@ -2582,25 +2698,31 @@ fn compactable_tombstone_ids(conn: &Connection) -> Result>>()?; let mut out = Vec::new(); for (op_id, origin, seq) in rows { + let local_vector: i64 = conn + .query_row( + "SELECT max_seq FROM sync_vectors + WHERE device_id = ?1", + [&origin], + |row| row.get(0), + ) + .optional()? + .unwrap_or(0); let mut all_acked = true; for device in &active_devices { - let ack: Option = conn - .query_row( + let seen = if device == &own_device_id { + local_vector + } else if device == &origin { + seq + } else { + conn.query_row( "SELECT max_seq FROM sync_peer_acks WHERE peer_device_id = ?1 AND origin_device_id = ?2", params![device, origin], |row| row.get(0), ) - .optional()?; - let local_vector: Option = conn - .query_row( - "SELECT max_seq FROM sync_vectors - WHERE device_id = ?1", - [&origin], - |row| row.get(0), - ) - .optional()?; - let seen = ack.or(local_vector).unwrap_or(0); + .optional()? + .unwrap_or(0) + }; if seen < seq { all_acked = false; break; @@ -2899,6 +3021,97 @@ mod tests { assert!(!device_revoked(&sync, device_id)); } + #[test] + fn tombstone_gc_waits_for_every_active_remote_ack() { + let sync = test_sync(); + let origin = sync.ensure_identity().unwrap().device_id; + sync.apply_device_trusted("dev_a", 1).unwrap(); + sync.apply_device_trusted("dev_b", 1).unwrap(); + + sync.record_local_op(SyncOpPayload::PlaylistDeleted { + playlist_id: "pl_deleted".to_string(), + }) + .unwrap(); + { + let conn = lock(&sync.conn); + let tombstones: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sync_ops WHERE tombstone = 1", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(tombstones, 1); + } + + let ack = BTreeMap::from([(origin, 1)]); + sync.note_peer_vector("dev_a", &ack).unwrap(); + sync.gc_tombstones().unwrap(); + { + let conn = lock(&sync.conn); + let tombstones: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sync_ops WHERE tombstone = 1", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(tombstones, 1); + } + + sync.note_peer_vector("dev_b", &ack).unwrap(); + sync.gc_tombstones().unwrap(); + let conn = lock(&sync.conn); + let tombstones: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sync_ops WHERE tombstone = 1", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(tombstones, 0); + } + + #[test] + fn snapshot_carries_deleted_playlists_to_repair_stale_peers() { + let source = test_sync(); + let source_playlist = source.library.create_playlist("Gone").unwrap(); + let playlist_sync_id = source + .library + .ensure_playlist_sync_id(source_playlist.id) + .unwrap(); + source + .apply_playlist_state(&playlist_sync_id, "Gone", false, 10, "dev_remote:1") + .unwrap(); + source + .apply_playlist_state(&playlist_sync_id, "", true, 20, "dev_remote:2") + .unwrap(); + let snapshot = source.snapshot().unwrap(); + assert!( + snapshot + .deleted_playlists + .iter() + .any(|playlist| playlist.playlist_id == playlist_sync_id) + ); + + let peer = test_sync(); + let peer_playlist = peer + .library + .upsert_synced_playlist(&playlist_sync_id, "Gone") + .unwrap(); + assert!(peer.library.playlist(peer_playlist).is_ok()); + + peer.apply_snapshot(snapshot).unwrap(); + assert!( + !peer + .library + .playlists() + .unwrap() + .iter() + .any(|playlist| playlist.title == "Gone") + ); + } + #[test] fn synced_fed_like_metadata_repairs_existing_like_state() { let sync = test_sync(); diff --git a/src/library/mod.rs b/src/library/mod.rs index 9da50e1..4bc6abb 100644 --- a/src/library/mod.rs +++ b/src/library/mod.rs @@ -867,6 +867,43 @@ impl Library { Ok(()) } + pub fn remove_content_ids_from_playlist( + &self, + playlist_id: i64, + content_ids: &[String], + ) -> Result<()> { + let conn = self.lock(); + let playlist_sync_id: Option = conn + .query_row( + "SELECT sync_id FROM playlists WHERE id = ?1", + [playlist_id], + |row| row.get::<_, Option>(0), + ) + .optional()? + .flatten(); + for content_id in content_ids { + let Some(content_id) = music_dht::normalize_content_id(content_id) else { + continue; + }; + conn.execute( + "DELETE FROM playlist_tracks + WHERE playlist_id = ?1 + AND track_id IN ( + SELECT id FROM tracks WHERE content_id = ?2 + )", + params![playlist_id, content_id], + )?; + if let Some(sync_id) = playlist_sync_id.as_deref() { + conn.execute( + "DELETE FROM fed_playlist_tracks + WHERE playlist_sync_id = ?1 AND content_id = ?2", + params![sync_id, content_id], + )?; + } + } + Ok(()) + } + pub fn track_content_id_by_id(&self, track_id: i64) -> Result> { let conn = self.lock(); let content_id: Option = conn @@ -881,26 +918,6 @@ impl Library { Ok(content_id) } - pub fn track_content_ids(&self, track_ids: &[i64]) -> Result> { - let conn = self.lock(); - let mut out = Vec::new(); - for &track_id in track_ids { - let content_id: Option = conn - .query_row( - "SELECT content_id FROM tracks WHERE id = ?1", - [track_id], - |row| row.get::<_, Option>(0), - ) - .optional()? - .flatten() - .and_then(|value| music_dht::normalize_content_id(&value)); - if let Some(content_id) = content_id { - out.push(content_id); - } - } - Ok(out) - } - pub fn playlist_track_content_positions( &self, playlist_id: i64, @@ -2115,9 +2132,14 @@ mod tests { .unwrap(); assert_eq!(card.track_count, 1); - lib.remove_content_id_from_synced_playlist(&sync_id, &content_id) + lib.remove_content_ids_from_playlist(playlist.id, std::slice::from_ref(&content_id)) .unwrap(); assert_eq!(lib.playlist(playlist.id).unwrap().tracks.len(), 0); + assert!( + lib.fed_playlist_track_by_content_id(&sync_id, &content_id) + .unwrap() + .is_none() + ); } #[test]