Connected Devices: fixed rplaylist sync

This commit is contained in:
Ultradesu
2026-07-24 02:35:36 +03:00
parent e08eb71ab7
commit 71d31505f9
3 changed files with 487 additions and 33 deletions
+6 -2
View File
@@ -1009,6 +1009,7 @@ pub(crate) fn fed_download_spawn(
tokio::spawn(async move {
let total = tracks.len();
let mut imported_ids = Vec::new();
let mut imported_fed_tracks = Vec::new();
let mut failed = 0usize;
for (index, track) in tracks.iter().enumerate() {
let _ = tx.send(AppEvent::StatusMessage(format!(
@@ -1017,7 +1018,10 @@ pub(crate) fn fed_download_spawn(
track.title
)));
match fed.download_to_library(track).await {
Ok(imported) => imported_ids.push(imported.id),
Ok(imported) => {
imported_ids.push(imported.id);
imported_fed_tracks.push(track.clone());
}
Err(err) => {
failed += 1;
tracing::warn!(title = %track.title, "federated download failed: {err:#}");
@@ -1041,7 +1045,7 @@ pub(crate) fn fed_download_spawn(
.map_err(|err| format!("{err:#}"));
if result.is_ok()
&& let Err(err) =
devices.record_playlist_tracks_added(playlist_id, &imported_ids)
devices.record_playlist_fed_tracks_added(playlist_id, &imported_fed_tracks)
{
tracing::warn!(%err, playlist_id, "recording synced playlist add failed");
}
+138 -20
View File
@@ -152,6 +152,8 @@ pub enum SyncOpPayload {
playlist_id: String,
content_id: String,
position: i64,
#[serde(default)]
fed: Option<SyncedFedTrack>,
},
PlaylistTrackRemoved {
playlist_id: String,
@@ -217,6 +219,8 @@ struct SnapshotPlaylistItem {
position: i64,
hlc_ms: i64,
op_id: String,
#[serde(default)]
fed: Option<SyncedFedTrack>,
}
enum PairAttempt {
@@ -722,17 +726,44 @@ impl DeviceSync {
}
pub fn record_playlist_tracks_added(&self, playlist_id: i64, track_ids: &[i64]) -> Result<()> {
let playlist_id = self.library.ensure_playlist_sync_id(playlist_id)?;
for (position, content_id) in self
let playlist_sync_id = self.library.ensure_playlist_sync_id(playlist_id)?;
for (content_id, position) in self
.library
.track_content_ids(track_ids)?
.into_iter()
.enumerate()
.playlist_track_content_positions(playlist_id, track_ids)?
{
self.record_local_op(SyncOpPayload::PlaylistTrackAdded {
playlist_id: playlist_id.clone(),
playlist_id: playlist_sync_id.clone(),
content_id,
position: position as i64,
position,
fed: None,
})?;
}
Ok(())
}
pub fn record_playlist_fed_tracks_added(
&self,
playlist_id: i64,
tracks: &[crate::federation::FedTrack],
) -> Result<()> {
let playlist_sync_id = self.library.ensure_playlist_sync_id(playlist_id)?;
for (fallback_position, fed) in tracks.iter().enumerate() {
let Some(content_id) = fed
.content_id
.as_deref()
.and_then(music_dht::normalize_content_id)
else {
continue;
};
let position = self
.library
.playlist_content_position(playlist_id, &content_id)?
.unwrap_or(fallback_position as i64);
self.record_local_op(SyncOpPayload::PlaylistTrackAdded {
playlist_id: playlist_sync_id.clone(),
content_id,
position,
fed: SyncedFedTrack::from_fed(fed),
})?;
}
Ok(())
@@ -1071,11 +1102,13 @@ impl DeviceSync {
playlist_id,
content_id,
position,
fed,
} => self.apply_playlist_item_state(
playlist_id,
content_id,
true,
*position,
fed.as_ref(),
op.hlc_ms,
&op.op_id,
)?,
@@ -1087,6 +1120,7 @@ impl DeviceSync {
content_id,
false,
0,
None,
op.hlc_ms,
&op.op_id,
)?,
@@ -1335,6 +1369,7 @@ impl DeviceSync {
content_id: &str,
present: bool,
position: i64,
fed: Option<&SyncedFedTrack>,
hlc_ms: i64,
op_id: &str,
) -> Result<bool> {
@@ -1377,14 +1412,23 @@ impl DeviceSync {
],
)?;
}
let mut visible_changed = apply;
if present {
self.library
visible_changed |= self
.library
.add_content_id_to_synced_playlist(playlist_id, content_id)?;
if let Some(fed) = fed {
visible_changed |= self.library.upsert_fed_playlist_track(
playlist_id,
&fed.to_fed_track(),
position,
)?;
}
} else {
self.library
.remove_content_id_from_synced_playlist(playlist_id, content_id)?;
}
Ok(true)
Ok(visible_changed)
}
fn apply_snapshot(&self, snapshot: SyncSnapshot) -> Result<()> {
@@ -1412,6 +1456,7 @@ impl DeviceSync {
&item.content_id,
true,
item.position,
item.fed.as_ref(),
item.hlc_ms,
&item.op_id,
)?;
@@ -1617,16 +1662,34 @@ impl DeviceSync {
WHERE playlist_id = ?1 AND present = 1
ORDER BY position, content_id",
)?;
let items = item_stmt
let item_rows = item_stmt
.query_map([&playlist_id], |row| {
Ok(SnapshotPlaylistItem {
content_id: row.get(0)?,
position: row.get(1)?,
hlc_ms: row.get(2)?,
op_id: row.get(3)?,
})
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, String>(3)?,
))
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
let mut items = Vec::new();
for (content_id, position, hlc_ms, op_id) in item_rows {
let fed_track = match self
.library
.fed_playlist_track_by_content_id(&playlist_id, &content_id)?
{
Some(fed) => Some(fed),
None => self.library.fed_like_by_content_id(&content_id)?,
};
let fed = fed_track.as_ref().and_then(SyncedFedTrack::from_fed);
items.push(SnapshotPlaylistItem {
content_id,
position,
hlc_ms,
op_id,
fed,
});
}
playlists.push(SnapshotPlaylist {
playlist_id,
title,
@@ -1808,17 +1871,22 @@ impl DeviceSync {
fn unresolved_playlist_item_count_with_conn(&self, conn: &Connection) -> Result<i64> {
let mut stmt = conn.prepare(
"SELECT content_id
"SELECT playlist_id, content_id
FROM sync_state_playlist_items
WHERE present = 1",
)?;
let ids = stmt
.query_map([], |row| row.get::<_, String>(0))?
.query_map([], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
drop(stmt);
let mut unresolved = 0;
for content_id in ids {
if self.library.track_id_by_content_id(&content_id)?.is_none() {
for (playlist_id, content_id) in ids {
if !self
.library
.has_playlist_content_reference(&playlist_id, &content_id)?
{
unresolved += 1;
}
}
@@ -2858,4 +2926,54 @@ mod tests {
);
assert!(sync.library.fed_like_ids().unwrap().is_empty());
}
#[test]
fn synced_playlist_item_metadata_creates_pending_fed_track() {
let sync = test_sync();
let playlist = sync.library.create_playlist("Remote Mix").unwrap();
let playlist_sync_id = sync.library.ensure_playlist_sync_id(playlist.id).unwrap();
let content_id = format!("b3:{}", "b".repeat(64));
let fed = test_fed_track(&content_id);
let synced = SyncedFedTrack::from_fed(&fed).unwrap();
assert!(
sync.apply_playlist_item_state(
&playlist_sync_id,
&content_id,
true,
3,
Some(&synced),
10,
"dev_remote:2",
)
.unwrap()
);
let detail = sync.library.playlist(playlist.id).unwrap();
assert_eq!(detail.tracks.len(), 1);
assert!(detail.tracks[0].is_fed_pending());
assert_eq!(detail.tracks[0].title, fed.title);
let conn = lock(&sync.conn);
assert_eq!(
sync.unresolved_playlist_item_count_with_conn(&conn)
.unwrap(),
0
);
drop(conn);
assert!(
sync.apply_playlist_item_state(
&playlist_sync_id,
&content_id,
false,
0,
None,
11,
"dev_remote:3",
)
.unwrap()
);
assert_eq!(sync.library.playlist(playlist.id).unwrap().tracks.len(), 0);
}
}
+343 -11
View File
@@ -10,6 +10,7 @@
pub mod import;
pub mod models;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Mutex;
@@ -97,6 +98,23 @@ CREATE TABLE IF NOT EXISTS fed_likes (
disc_number INTEGER,
liked_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE TABLE IF NOT EXISTS fed_playlist_tracks (
playlist_sync_id TEXT NOT NULL,
item_id TEXT NOT NULL,
owner TEXT NOT NULL,
title TEXT NOT NULL,
artist_names TEXT NOT NULL DEFAULT '',
featured_artist_names TEXT NOT NULL DEFAULT '',
year INTEGER,
duration_seconds REAL,
content_id TEXT NOT NULL,
release_title TEXT,
track_number INTEGER,
disc_number INTEGER,
position INTEGER NOT NULL DEFAULT 0,
added_at TEXT NOT NULL DEFAULT (datetime('now')),
PRIMARY KEY (playlist_sync_id, content_id)
);
CREATE TABLE IF NOT EXISTS history (
id INTEGER PRIMARY KEY,
track_id INTEGER NOT NULL REFERENCES tracks(id) ON DELETE CASCADE,
@@ -110,6 +128,8 @@ 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_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);
";
/// The SELECT column list every TrackItem row is built from; artist lists
@@ -599,6 +619,15 @@ impl Library {
let mut statement = conn.prepare(
"SELECT p.id, p.title,
(SELECT COUNT(*) FROM playlist_tracks pt WHERE pt.playlist_id = p.id)
+
(SELECT COUNT(*) FROM fed_playlist_tracks f
WHERE f.playlist_sync_id = p.sync_id
AND NOT EXISTS (
SELECT 1 FROM playlist_tracks pt
JOIN tracks t ON t.id = pt.track_id
WHERE pt.playlist_id = p.id
AND t.content_id = f.content_id
))
FROM playlists p ORDER BY p.title COLLATE NOCASE",
)?;
let rows = statement
@@ -641,14 +670,14 @@ impl Library {
tracks,
});
}
let (title, description): (String, Option<String>) = conn
let (title, description, sync_id): (String, Option<String>, Option<String>) = conn
.query_row(
"SELECT title, description FROM playlists WHERE id = ?1",
"SELECT title, description, sync_id FROM playlists WHERE id = ?1",
[id],
|row| Ok((row.get(0)?, row.get(1)?)),
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
)
.context("playlist not found")?;
let tracks = query_tracks(
let local_tracks = query_tracks(
&conn,
&format!(
"SELECT {TRACK_COLUMNS} FROM tracks t
@@ -659,11 +688,40 @@ impl Library {
),
params![id],
)?;
let mut positions = HashMap::new();
let mut position_stmt =
conn.prepare("SELECT track_id, position FROM playlist_tracks WHERE playlist_id = ?1")?;
let position_rows = position_stmt.query_map([id], |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?))
})?;
for row in position_rows {
let (track_id, position) = row?;
positions.insert(track_id, position);
}
drop(position_stmt);
let mut entries: Vec<(i64, TrackItem)> = local_tracks
.into_iter()
.enumerate()
.map(|(index, track)| {
let position = positions.get(&track.id).copied().unwrap_or(index as i64);
(position, track)
})
.collect();
if let Some(sync_id) = sync_id.as_deref() {
for (position, fed) in fed_playlist_tracks(&conn, id, sync_id)? {
entries.push((position, crate::federation::pending_track(&fed)));
}
}
entries.sort_by(|(left_pos, left_track), (right_pos, right_track)| {
left_pos
.cmp(right_pos)
.then_with(|| left_track.title.cmp(&right_track.title))
});
Ok(PlaylistDetail {
id,
title,
description,
tracks,
tracks: entries.into_iter().map(|(_, track)| track).collect(),
})
}
@@ -695,12 +753,28 @@ impl Library {
pub fn delete_playlist(&self, id: i64) -> Result<()> {
let conn = self.lock();
let sync_id: Option<String> = conn
.query_row("SELECT sync_id FROM playlists WHERE id = ?1", [id], |row| {
row.get::<_, Option<String>>(0)
})
.optional()?
.flatten();
if let Some(sync_id) = sync_id {
conn.execute(
"DELETE FROM fed_playlist_tracks WHERE playlist_sync_id = ?1",
[sync_id],
)?;
}
conn.execute("DELETE FROM playlists WHERE id = ?1", [id])?;
Ok(())
}
pub fn delete_playlist_by_sync_id(&self, sync_id: &str) -> Result<()> {
let conn = self.lock();
conn.execute(
"DELETE FROM fed_playlist_tracks WHERE playlist_sync_id = ?1",
[sync_id],
)?;
conn.execute("DELETE FROM playlists WHERE sync_id = ?1", [sync_id])?;
Ok(())
}
@@ -827,6 +901,59 @@ impl Library {
Ok(out)
}
pub fn playlist_track_content_positions(
&self,
playlist_id: i64,
track_ids: &[i64],
) -> Result<Vec<(String, i64)>> {
let conn = self.lock();
let mut out = Vec::new();
for &track_id in track_ids {
let row: Option<(Option<String>, i64)> = conn
.query_row(
"SELECT t.content_id, pt.position
FROM tracks t
JOIN playlist_tracks pt ON pt.track_id = t.id
WHERE t.id = ?1 AND pt.playlist_id = ?2",
params![track_id, playlist_id],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.optional()?;
let Some((content_id, position)) = row else {
continue;
};
if let Some(content_id) = content_id
.as_deref()
.and_then(music_dht::normalize_content_id)
{
out.push((content_id, position));
}
}
Ok(out)
}
pub fn playlist_content_position(
&self,
playlist_id: i64,
content_id: &str,
) -> Result<Option<i64>> {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(None);
};
let conn = self.lock();
Ok(conn
.query_row(
"SELECT pt.position
FROM playlist_tracks pt
JOIN tracks t ON t.id = pt.track_id
WHERE pt.playlist_id = ?1 AND t.content_id = ?2
LIMIT 1",
params![playlist_id, content_id],
|row| row.get(0),
)
.optional()?)
}
pub fn track_id_by_content_id(&self, content_id: &str) -> Result<Option<i64>> {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(None);
@@ -854,6 +981,91 @@ impl Library {
Ok(changed > 0)
}
pub fn upsert_fed_playlist_track(
&self,
playlist_sync_id: &str,
fed: &crate::federation::FedTrack,
position: i64,
) -> Result<bool> {
let Some(content_id) = fed
.content_id
.as_deref()
.and_then(music_dht::normalize_content_id)
else {
return Ok(false);
};
let conn = self.lock();
Ok(conn.execute(
"INSERT INTO fed_playlist_tracks
(playlist_sync_id, item_id, owner, title, artist_names,
featured_artist_names, year, duration_seconds, content_id,
release_title, track_number, disc_number, position)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)
ON CONFLICT(playlist_sync_id, content_id) DO UPDATE SET
item_id = excluded.item_id,
owner = excluded.owner,
title = excluded.title,
artist_names = excluded.artist_names,
featured_artist_names = excluded.featured_artist_names,
year = excluded.year,
duration_seconds = excluded.duration_seconds,
release_title = excluded.release_title,
track_number = excluded.track_number,
disc_number = excluded.disc_number,
position = excluded.position",
params![
playlist_sync_id,
fed.item_id,
fed.owner,
fed.title,
fed.artist_names.join("; "),
fed.featured_artist_names.join("; "),
fed.year,
fed.duration_seconds.map(|d| d as f64),
content_id,
fed.release_title,
fed.track_number,
fed.disc_number,
position,
],
)? > 0)
}
pub fn fed_playlist_track_by_content_id(
&self,
playlist_sync_id: &str,
content_id: &str,
) -> Result<Option<crate::federation::FedTrack>> {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(None);
};
let conn = self.lock();
conn.query_row(
"SELECT item_id, owner, title, artist_names, featured_artist_names,
year, duration_seconds, content_id, release_title, track_number, disc_number
FROM fed_playlist_tracks
WHERE playlist_sync_id = ?1 AND content_id = ?2
LIMIT 1",
params![playlist_sync_id, content_id],
fed_track_from_row,
)
.optional()
.map_err(Into::into)
}
pub fn has_playlist_content_reference(
&self,
playlist_sync_id: &str,
content_id: &str,
) -> Result<bool> {
if self.track_id_by_content_id(content_id)?.is_some() {
return Ok(true);
}
Ok(self
.fed_playlist_track_by_content_id(playlist_sync_id, content_id)?
.is_some())
}
pub fn add_content_id_to_synced_playlist(
&self,
playlist_sync_id: &str,
@@ -891,7 +1103,7 @@ impl Library {
playlist_sync_id: &str,
content_id: &str,
) -> Result<()> {
let Some(track_id) = self.track_id_by_content_id(content_id)? else {
let Some(content_id) = music_dht::normalize_content_id(content_id) else {
return Ok(());
};
let conn = self.lock();
@@ -902,12 +1114,25 @@ impl Library {
|row| row.get::<_, i64>(0),
)
.optional()?;
let Some(playlist_id) = playlist_id else {
return Ok(());
};
if let Some(playlist_id) = playlist_id {
let track_id: Option<i64> = conn
.query_row(
"SELECT id FROM tracks WHERE content_id = ?1 LIMIT 1",
[&content_id],
|row| row.get(0),
)
.optional()?;
if let Some(track_id) = track_id {
conn.execute(
"DELETE FROM playlist_tracks WHERE playlist_id = ?1 AND track_id = ?2",
params![playlist_id, track_id],
)?;
}
}
conn.execute(
"DELETE FROM playlist_tracks WHERE playlist_id = ?1 AND track_id = ?2",
params![playlist_id, track_id],
"DELETE FROM fed_playlist_tracks
WHERE playlist_sync_id = ?1 AND content_id = ?2",
params![playlist_sync_id, content_id],
)?;
Ok(())
}
@@ -1539,6 +1764,64 @@ fn query_tracks(
Ok(tracks)
}
fn fed_playlist_tracks(
conn: &Connection,
playlist_id: i64,
playlist_sync_id: &str,
) -> Result<Vec<(i64, crate::federation::FedTrack)>> {
let mut statement = conn.prepare(
"SELECT position, item_id, owner, title, artist_names, featured_artist_names,
year, duration_seconds, content_id, release_title, track_number, disc_number
FROM fed_playlist_tracks f
WHERE f.playlist_sync_id = ?1
AND NOT EXISTS (
SELECT 1 FROM playlist_tracks pt
JOIN tracks t ON t.id = pt.track_id
WHERE pt.playlist_id = ?2
AND t.content_id = f.content_id
)
ORDER BY position, title COLLATE NOCASE",
)?;
let rows = statement.query_map(params![playlist_sync_id, playlist_id], |row| {
Ok((row.get::<_, i64>(0)?, fed_track_from_offset_row(row, 1)?))
})?;
Ok(rows.collect::<rusqlite::Result<Vec<_>>>()?)
}
fn fed_track_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<crate::federation::FedTrack> {
fed_track_from_offset_row(row, 0)
}
fn fed_track_from_offset_row(
row: &rusqlite::Row<'_>,
offset: usize,
) -> rusqlite::Result<crate::federation::FedTrack> {
Ok(crate::federation::FedTrack {
item_id: row.get(offset)?,
owner: row.get(offset + 1)?,
own: false,
title: row.get(offset + 2)?,
artist_names: split_joined_names(row.get(offset + 3)?),
featured_artist_names: split_joined_names(row.get(offset + 4)?),
year: row.get(offset + 5)?,
duration_seconds: row
.get::<_, Option<f64>>(offset + 6)?
.map(|d| d.round() as i64),
content_id: row.get(offset + 7)?,
release_title: row.get(offset + 8)?,
track_number: row.get(offset + 9)?,
disc_number: row.get(offset + 10)?,
})
}
fn split_joined_names(names: String) -> Vec<String> {
names
.split("; ")
.filter(|name| !name.is_empty())
.map(str::to_string)
.collect()
}
/// Escape LIKE wildcards in user input; queries use `ESCAPE '\'`.
/// Registers `norm(text)` — Unicode-aware case folding and normalization
/// (NFKC, lowercase, punctuation stripped). SQLite's own LIKE/NOCASE only
@@ -1788,6 +2071,55 @@ mod tests {
assert_eq!(lib.playlists().unwrap().len(), 1);
}
#[test]
fn synced_playlist_can_show_federated_pending_tracks() {
let lib = test_library();
let playlist = lib.create_playlist("Remote Mix").unwrap();
let sync_id = lib.ensure_playlist_sync_id(playlist.id).unwrap();
let content_id = format!("b3:{}", "a".repeat(64));
let fed = crate::federation::FedTrack {
item_id: "fed_item_1".to_string(),
owner: "fed_owner_1".to_string(),
own: false,
title: "Remote Song".to_string(),
artist_names: vec!["Remote Artist".to_string()],
featured_artist_names: vec!["Remote Guest".to_string()],
year: Some(2026),
duration_seconds: Some(123),
content_id: Some(content_id.clone()),
release_title: Some("Remote Release".to_string()),
track_number: Some(2),
disc_number: Some(1),
};
assert!(lib.upsert_fed_playlist_track(&sync_id, &fed, 4).unwrap());
assert!(
lib.has_playlist_content_reference(&sync_id, &content_id)
.unwrap()
);
let detail = lib.playlist(playlist.id).unwrap();
assert_eq!(detail.tracks.len(), 1);
let track = &detail.tracks[0];
assert!(track.is_fed_pending());
assert_eq!(track.title, "Remote Song");
assert_eq!(track.artist_line(), "Remote Artist feat. Remote Guest");
assert_eq!(track.release_title, "Remote Release");
assert_eq!(track.content_id.as_deref(), Some(content_id.as_str()));
let card = lib
.playlists()
.unwrap()
.into_iter()
.find(|card| card.id == playlist.id)
.unwrap();
assert_eq!(card.track_count, 1);
lib.remove_content_id_from_synced_playlist(&sync_id, &content_id)
.unwrap();
assert_eq!(lib.playlist(playlist.id).unwrap().tracks.len(), 0);
}
#[test]
fn track_edit_relinks_artists() {
let lib = test_library();