Connected Devices: fixed rplaylist sync

This commit is contained in:
Ultradesu
2026-07-24 03:06:20 +03:00
parent 71d31505f9
commit 682794755e
4 changed files with 291 additions and 40 deletions
+7 -2
View File
@@ -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");
}
+11
View File
@@ -37,6 +37,7 @@ pub enum Effect {
RemoveFromPlaylist {
playlist_id: i64,
track_ids: Vec<i64>,
content_ids: Vec<String>,
},
RemoveQueueIndices {
indices: Vec<usize>,
@@ -607,6 +608,15 @@ fn delete_selected(state: &mut AppState) -> Option<Effect> {
return None;
}
let track_ids: Vec<i64> = tracks.iter().map(|track| track.id).collect();
let content_ids: Vec<String> = 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<i64> = track_ids
@@ -626,6 +636,7 @@ fn delete_selected(state: &mut AppState) -> Option<Effect> {
return Some(Effect::RemoveFromPlaylist {
playlist_id: opened.id,
track_ids,
content_ids,
});
}
}
+230 -17
View File
@@ -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<SnapshotLike>,
#[serde(default)]
unlikes: Vec<SnapshotLikeTombstone>,
#[serde(default)]
playlists: Vec<SnapshotPlaylist>,
#[serde(default)]
deleted_playlists: Vec<SnapshotPlaylistTombstone>,
#[serde(default)]
removed_playlist_items: Vec<SnapshotPlaylistItemTombstone>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -203,6 +209,13 @@ struct SnapshotLike {
fed: Option<SyncedFedTrack>,
}
#[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<SnapshotPlaylistItem>,
}
#[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<SyncedFedTrack>,
}
#[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::<rusqlite::Result<Vec<_>>>()?
};
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::<rusqlite::Result<Vec<_>>>()?
};
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::<rusqlite::Result<Vec<_>>>()?
};
Ok(SyncSnapshot {
likes,
unlikes,
playlists,
deleted_playlists,
removed_playlist_items,
})
}
fn device_profiles(&self) -> Result<Vec<DeviceProfileWire>> {
@@ -2548,6 +2663,7 @@ fn payload_kind(payload: &SyncOpPayload) -> &'static str {
}
fn compactable_tombstone_ids(conn: &Connection) -> Result<Vec<(String, String, i64)>> {
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<Vec<(String, String, i
.collect::<rusqlite::Result<Vec<_>>>()?;
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<i64> = 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<i64> = 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();
+43 -21
View File
@@ -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<String> = conn
.query_row(
"SELECT sync_id FROM playlists WHERE id = ?1",
[playlist_id],
|row| row.get::<_, Option<String>>(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<Option<String>> {
let conn = self.lock();
let content_id: Option<String> = conn
@@ -881,26 +918,6 @@ impl Library {
Ok(content_id)
}
pub fn track_content_ids(&self, track_ids: &[i64]) -> Result<Vec<String>> {
let conn = self.lock();
let mut out = Vec::new();
for &track_id in track_ids {
let content_id: Option<String> = conn
.query_row(
"SELECT content_id FROM tracks WHERE id = ?1",
[track_id],
|row| row.get::<_, Option<String>>(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]