1use std::collections::{HashMap, HashSet};
2use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
3
4use crate::db::connection::Database;
5use crate::db::queries::{self, TrackMeta};
6use crate::remote::client::{
7 SubsonicAlbum, SubsonicAlbumFull, SubsonicArtist, SubsonicClient, SubsonicError, SubsonicSong,
8};
9
10use rayon::prelude::*;
11use rusqlite::params;
12use thiserror::Error;
13
14#[derive(Debug, Error)]
15pub enum SyncError {
16 #[error("subsonic error: {0}")]
17 Subsonic(#[from] super::client::SubsonicError),
18 #[error("db error: {0}")]
19 Db(#[from] crate::db::connection::DbError),
20}
21
22#[derive(Debug, Default)]
23pub struct SyncResult {
24 pub artists_synced: usize,
25 pub albums_synced: usize,
26 pub tracks_synced: usize,
27 pub albums_failed: usize,
30 pub pages_failed: usize,
33 pub tracks_failed: usize,
36 pub tracks_removed: usize,
38}
39
40impl SyncResult {
41 pub fn is_complete(&self) -> bool {
43 self.albums_failed == 0 && self.pages_failed == 0 && self.tracks_failed == 0
44 }
45}
46
47pub fn get_last_sync(
49 db: &Database,
50 url: &str,
51) -> Result<Option<i64>, crate::db::connection::DbError> {
52 let result = db.conn.query_row(
53 "SELECT last_sync FROM remote_servers WHERE url = ?1",
54 params![url],
55 |row| row.get::<_, Option<i64>>(0),
56 );
57 match result {
58 Ok(ts) => Ok(ts),
59 Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
60 Err(e) => Err(e.into()),
61 }
62}
63
64pub fn library_version(db: &Database, url: &str) -> Option<i64> {
67 db.conn
68 .query_row(
69 "SELECT library_version FROM remote_servers WHERE url = ?1",
70 params![url],
71 |row| row.get::<_, Option<i64>>(0),
72 )
73 .ok()
74 .flatten()
75}
76
77pub fn set_library_version(
79 db: &Database,
80 url: &str,
81 username: &str,
82 version: i64,
83) -> Result<(), crate::db::connection::DbError> {
84 db.conn.execute(
85 "INSERT INTO remote_servers (url, username, library_version)
86 VALUES (?1, ?2, ?3)
87 ON CONFLICT(url) DO UPDATE SET library_version = ?3",
88 params![url, username, version],
89 )?;
90 Ok(())
91}
92
93pub fn update_last_sync(
95 db: &Database,
96 url: &str,
97 username: &str,
98 timestamp: i64,
99) -> Result<(), crate::db::connection::DbError> {
100 db.conn.execute(
101 "INSERT INTO remote_servers (url, username, last_sync)
102 VALUES (?1, ?2, ?3)
103 ON CONFLICT(url) DO UPDATE SET last_sync = ?3",
104 params![url, username, timestamp],
105 )?;
106 Ok(())
107}
108
109#[derive(Debug, Clone, Copy, PartialEq, Eq)]
111pub enum SyncPhase {
112 Albums,
114 Tracks,
116 Artists,
118 Finishing,
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126pub struct SyncProgress {
127 pub phase: SyncPhase,
128 pub done: u64,
129 pub total: Option<u64>,
130}
131
132const PAGE_SIZE: u32 = 500;
134
135const FETCH_LANES: usize = 4;
139
140const PAGE_ATTEMPTS: u32 = 3;
142
143pub fn sync_library(
164 db: &Database,
165 client: &SubsonicClient,
166 server_url: &str,
167 username: &str,
168 progress: &(dyn Fn(SyncProgress) + Sync),
169) -> Result<SyncResult, SyncError> {
170 let mut result = SyncResult::default();
171 log::info!("syncing the library from {server_url}");
172
173 let sync_start = std::time::SystemTime::now()
174 .duration_since(std::time::UNIX_EPOCH)
175 .unwrap_or_default()
176 .as_secs() as i64;
177
178 let albums = list_albums(client, progress)?;
179
180 let expected: u64 = albums
181 .iter()
182 .filter_map(|a| a.song_count)
183 .map(|n| n.max(0) as u64)
184 .sum();
185
186 let mut song_ids: HashSet<String> = HashSet::new();
187 let mut total = (expected > 0).then_some(expected);
188 if total.is_none() {
189 total = client.song_count().ok().flatten();
190 }
191 let walked = !albums.is_empty()
192 && sync_all_songs(
193 db,
194 client,
195 &albums,
196 total,
197 &mut result,
198 &mut song_ids,
199 progress,
200 )?;
201 if !walked {
202 let all: Vec<&SubsonicAlbum> = albums.iter().collect();
203 sync_by_album(
204 db,
205 client,
206 &all,
207 total,
208 &mut result,
209 &mut song_ids,
210 progress,
211 )?;
212 }
213
214 let artists = client.get_artists()?;
218 result.artists_synced = artists.len();
219 progress(SyncProgress {
220 phase: SyncPhase::Artists,
221 done: 0,
222 total: Some(artists.len() as u64),
223 });
224 write_artists(db, &artists, &mut result);
225
226 progress(SyncProgress {
227 phase: SyncPhase::Finishing,
228 done: 0,
229 total: None,
230 });
231
232 let mut live_albums: std::collections::HashSet<String> =
238 albums.iter().map(|a| a.id.clone()).collect();
239 if result.is_complete() && !live_albums.is_empty() {
240 confirm_missing_albums(db, client, &mut live_albums);
241 }
242 let albums_short = live_albums.len() > albums.len();
243
244 let listed_everything = !albums_short && total.is_none_or(|n| song_ids.len() as u64 >= n);
249 if result.is_complete() && !song_ids.is_empty() && listed_everything {
250 match queries::relink_vanished_remote_ids(&db.conn, &song_ids) {
251 Ok(0) => {}
252 Ok(n) => log::info!("{n} files had ids the server no longer knows; relinked"),
253 Err(e) => log::warn!("failed to relink tracks with vanished remote ids: {e}"),
254 }
255 }
256
257 let live_tracks =
261 (result.is_complete() && !song_ids.is_empty() && listed_everything).then_some(&song_ids);
262 if result.is_complete() && !live_albums.is_empty() {
263 match queries::remove_vanished_remote(&db.conn, live_tracks, Some(&live_albums)) {
264 Ok(0) => {}
265 Ok(n) => {
266 result.tracks_removed = n;
267 log::info!("{n} tracks the server no longer has; removed");
268 }
269 Err(e) => log::warn!("failed to remove tracks the server deleted: {e}"),
270 }
271 }
272
273 if result.is_complete() {
274 update_last_sync(db, server_url, username, sync_start)?;
275 } else {
276 log::warn!(
277 "{} album(s) and {} page(s) failed to fetch and {} track(s) failed to write — the next sync walks the library again",
278 result.albums_failed,
279 result.pages_failed,
280 result.tracks_failed,
281 );
282 }
283
284 log::info!(
285 "sync complete: {} artists, {} albums, {} tracks, {} albums and {} tracks failed",
286 result.artists_synced,
287 result.albums_synced,
288 result.tracks_synced,
289 result.albums_failed,
290 result.tracks_failed,
291 );
292
293 db.optimize();
294
295 Ok(result)
296}
297
298fn confirm_missing_albums(db: &Database, client: &SubsonicClient, live: &mut HashSet<String>) {
304 let held: Vec<String> = db
305 .conn
306 .prepare(
307 "SELECT DISTINCT al.remote_id FROM albums al JOIN tracks t ON t.album_id = al.id
308 WHERE al.remote_id IS NOT NULL AND t.path IS NULL AND t.remote_id IS NOT NULL",
309 )
310 .and_then(|mut stmt| {
311 stmt.query_map([], |r| r.get(0))?
312 .collect::<rusqlite::Result<Vec<String>>>()
313 })
314 .unwrap_or_default();
315 let missing: Vec<String> = held.into_iter().filter(|id| !live.contains(id)).collect();
316 for id in missing {
317 match client.get_album(&id) {
318 Err(SubsonicError::Api { code: 70, .. }) => {}
319 Ok(_) => {
320 log::info!("album {id} was missing from the listing but still exists");
321 live.insert(id);
322 }
323 Err(e) => {
324 log::warn!("could not confirm album {id} was deleted ({e}); keeping it");
325 live.insert(id);
326 }
327 }
328 }
329}
330
331fn list_albums(
333 client: &SubsonicClient,
334 progress: &(dyn Fn(SyncProgress) + Sync),
335) -> Result<Vec<SubsonicAlbum>, SyncError> {
336 let mut albums = Vec::new();
337 let mut seen: HashSet<String> = HashSet::new();
340 let mut offset = 0u32;
341 loop {
342 let page = client.get_album_list("alphabeticalByName", PAGE_SIZE, offset)?;
343 let count = page.len() as u32;
344 offset += count;
345 albums.extend(page.into_iter().filter(|a| seen.insert(a.id.clone())));
346 progress(SyncProgress {
347 phase: SyncPhase::Albums,
348 done: albums.len() as u64,
349 total: None,
350 });
351 if count < PAGE_SIZE {
352 return Ok(albums);
353 }
354 }
355}
356
357fn sync_all_songs(
364 db: &Database,
365 client: &SubsonicClient,
366 albums: &[SubsonicAlbum],
367 total: Option<u64>,
368 result: &mut SyncResult,
369 song_ids: &mut HashSet<String>,
370 progress: &(dyn Fn(SyncProgress) + Sync),
371) -> Result<bool, SyncError> {
372 let first = match client.all_songs_page(PAGE_SIZE, 0) {
373 Ok(songs) if !songs.is_empty() => songs,
374 Ok(_) => {
375 log::info!("server lists no songs for an empty search; syncing album by album");
376 return Ok(false);
377 }
378 Err(e) => {
379 log::info!("server refused an empty search ({e}); syncing album by album");
380 return Ok(false);
381 }
382 };
383
384 let by_id: HashMap<&str, &SubsonicAlbum> = albums.iter().map(|a| (a.id.as_str(), a)).collect();
385 let mut albums_seen: HashSet<String> = HashSet::new();
386 let mut write = |songs: Vec<SubsonicSong>,
387 result: &mut SyncResult,
388 song_ids: &mut HashSet<String>|
389 -> Result<(), SyncError> {
390 let batch = group_by_album(songs, &by_id);
391 albums_seen.extend(batch.iter().map(|a| a.id.clone()));
392 write_albums(db, client, &batch, result, song_ids)?;
393 progress(SyncProgress {
394 phase: SyncPhase::Tracks,
395 done: song_ids.len() as u64,
396 total,
397 });
398 Ok(())
399 };
400
401 let short = (first.len() as u32) < PAGE_SIZE;
402 write(first, result, song_ids)?;
403
404 if !short {
405 let next = AtomicU32::new(PAGE_SIZE);
411 let end = AtomicU32::new(u32::MAX);
414 let (tx, rx) = std::sync::mpsc::sync_channel(FETCH_LANES);
415 let failed = std::thread::scope(|scope| -> Result<usize, SyncError> {
416 for _ in 0..FETCH_LANES {
417 let tx = tx.clone();
418 let (next, end) = (&next, &end);
419 scope.spawn(move || {
420 loop {
421 let offset = next.fetch_add(PAGE_SIZE, Ordering::Relaxed);
422 if offset >= end.load(Ordering::Relaxed) {
423 return;
424 }
425 let page = fetch_page(client, offset);
426 if page
431 .as_ref()
432 .map_or(true, |songs| (songs.len() as u32) < PAGE_SIZE)
433 {
434 end.fetch_min(offset, Ordering::Relaxed);
435 }
436 if tx.send((offset, page)).is_err() {
437 return;
438 }
439 }
440 });
441 }
442 drop(tx);
443
444 let mut failed = 0;
445 for (offset, page) in rx {
446 match page {
447 Ok(songs) => write(songs, result, song_ids)?,
448 Err(e) => {
449 log::warn!("failed to fetch songs from offset {offset}: {e}");
450 failed += 1;
451 }
452 }
453 }
454 Ok(failed)
455 })?;
456 result.pages_failed += failed;
457 }
458
459 result.albums_synced += albums_seen.len();
460 Ok(true)
461}
462
463fn fetch_page(
466 client: &SubsonicClient,
467 offset: u32,
468) -> Result<Vec<SubsonicSong>, super::client::SubsonicError> {
469 let mut attempt = 1;
470 loop {
471 match client.all_songs_page(PAGE_SIZE, offset) {
472 Ok(songs) => return Ok(songs),
473 Err(e) if attempt >= PAGE_ATTEMPTS => return Err(e),
474 Err(e) => {
475 log::debug!("songs from offset {offset}, attempt {attempt}: {e}");
476 std::thread::sleep(std::time::Duration::from_millis(250 * attempt as u64));
477 attempt += 1;
478 }
479 }
480 }
481}
482
483fn group_by_album(
489 songs: Vec<SubsonicSong>,
490 albums: &HashMap<&str, &SubsonicAlbum>,
491) -> Vec<SubsonicAlbumFull> {
492 let mut grouped: Vec<SubsonicAlbumFull> = Vec::new();
493 let mut index: HashMap<String, usize> = HashMap::new();
494 for song in songs {
495 let Some(album_id) = song.album_id.clone() else {
496 log::warn!("song {} has no album id; skipped", song.id);
497 continue;
498 };
499 let i = *index.entry(album_id.clone()).or_insert_with(|| {
500 grouped.push(match albums.get(album_id.as_str()) {
501 Some(album) => SubsonicAlbumFull {
502 id: album.id.clone(),
503 name: album.name.clone(),
504 artist: album.artist.clone(),
505 artist_id: album.artist_id.clone(),
506 year: album.year,
507 genre: album.genre.clone(),
508 song_count: album.song_count,
509 created: album.created.clone(),
510 music_brainz_id: album.music_brainz_id.clone(),
511 sort_name: album.sort_name.clone(),
512 record_labels: album.record_labels.clone(),
513 song: Vec::new(),
514 },
515 None => SubsonicAlbumFull {
516 id: album_id,
517 name: song.album.clone().unwrap_or_default(),
518 artist: song.artist.clone(),
519 artist_id: song.artist_id.clone(),
520 year: song.year,
521 genre: song.genre.clone(),
522 song_count: None,
523 created: None,
524 music_brainz_id: None,
525 sort_name: None,
526 record_labels: Vec::new(),
527 song: Vec::new(),
528 },
529 });
530 grouped.len() - 1
531 });
532 grouped[i].song.push(song);
533 }
534 grouped
535}
536
537const ALBUM_BATCH: usize = 100;
540
541fn sync_by_album(
544 db: &Database,
545 client: &SubsonicClient,
546 albums: &[&SubsonicAlbum],
547 total: Option<u64>,
548 result: &mut SyncResult,
549 song_ids: &mut HashSet<String>,
550 progress: &(dyn Fn(SyncProgress) + Sync),
551) -> Result<(), SyncError> {
552 for batch in albums.chunks(ALBUM_BATCH) {
553 let failures = AtomicUsize::new(0);
554 let fetched: Vec<SubsonicAlbumFull> = batch
555 .par_iter()
556 .filter_map(|album| match client.get_album(&album.id) {
557 Ok(full) => Some(full),
558 Err(e) => {
559 log::warn!("failed to fetch album {}: {}", album.id, e);
560 failures.fetch_add(1, Ordering::Relaxed);
561 None
562 }
563 })
564 .collect();
565 result.albums_failed += failures.into_inner();
566 result.albums_synced += fetched.len();
567
568 write_albums(db, client, &fetched, result, song_ids)?;
569 progress(SyncProgress {
570 phase: SyncPhase::Tracks,
571 done: song_ids.len() as u64,
572 total,
573 });
574 log::info!(
575 "synced {} albums ({} tracks) so far...",
576 result.albums_synced,
577 result.tracks_synced
578 );
579 }
580 Ok(())
581}
582
583fn write_artists(db: &Database, artists: &[SubsonicArtist], result: &mut SyncResult) {
588 if db.conn.execute_batch("BEGIN IMMEDIATE").is_err() {
592 return;
593 }
594 let mut enriched = 0;
595 for artist in artists {
596 match queries::enrich_remote_artist(
597 &db.conn,
598 &artist.id,
599 artist.music_brainz_id.as_deref(),
600 artist.sort_name.as_deref(),
601 ) {
602 Ok(()) => enriched += 1,
603 Err(e) => log::warn!(
604 "failed to record artist metadata for {}: {}",
605 artist.name,
606 e
607 ),
608 }
609 }
610 if db.conn.execute_batch("COMMIT").is_err() {
611 let _ = db.conn.execute_batch("ROLLBACK");
612 return;
613 }
614 result.artists_synced = enriched;
615 log::info!("recorded metadata for {enriched} artists");
616}
617
618fn write_albums(
620 db: &Database,
621 client: &SubsonicClient,
622 albums: &[SubsonicAlbumFull],
623 result: &mut SyncResult,
624 song_ids: &mut HashSet<String>,
625) -> Result<(), SyncError> {
626 db.conn
629 .execute_batch("BEGIN IMMEDIATE")
630 .map_err(crate::db::connection::DbError::from)?;
631
632 for album in albums {
633 let artist_name = album.artist.as_deref().unwrap_or("Unknown Artist");
634
635 for song in &album.song {
636 song_ids.insert(song.id.clone());
637 let meta = TrackMeta {
638 title: song.title.clone(),
639 artist: song
640 .artist
641 .clone()
642 .unwrap_or_else(|| artist_name.to_string()),
643 album_artist: album.artist.clone(),
644 album: album.name.clone(),
645 date: album.year.map(|y| y.to_string()),
646 disc: song.disc_number,
647 track_number: song.track,
648 genre: song.genre.clone().or_else(|| album.genre.clone()),
649 label: None,
650 duration_ms: song.duration.map(|d| d * 1000),
651 codec: song.suffix.clone(),
652 sample_rate: positive(song.sampling_rate),
659 bit_depth: positive(song.bit_depth),
660 channels: positive(song.channel_count),
661 bitrate: song.bit_rate,
662 size_bytes: None,
663 mtime: None,
664 path: None,
665 source: "remote".to_string(),
666 remote_id: Some(song.id.clone()),
667 remote_url: Some(client.stream_url_template(&song.id)),
668 album_remote_id: Some(album.id.clone()),
669 artist_remote_id: album.artist_id.clone(),
670 mbid: song.music_brainz_id.clone(),
671 album_mbid: album.music_brainz_id.clone(),
672 album_added_at: album.created.clone(),
673 };
674
675 match queries::upsert_synced_track(&db.conn, &meta, song_ids) {
676 Ok(_) => result.tracks_synced += 1,
677 Err(e) => {
678 result.tracks_failed += 1;
679 log::warn!("failed to insert remote track {}: {}", song.title, e);
680 }
681 }
682 }
683
684 if let Err(e) = queries::enrich_remote_album(
688 &db.conn,
689 &album.id,
690 album.music_brainz_id.as_deref(),
691 album.sort_name.as_deref(),
692 album.song_count,
693 album
694 .record_labels
695 .iter()
696 .map(|l| l.name.as_str())
697 .find(|l| !l.is_empty()),
698 ) {
699 log::warn!("failed to record album metadata for {}: {}", album.name, e);
700 }
701 }
702
703 db.conn
704 .execute_batch("COMMIT")
705 .map_err(crate::db::connection::DbError::from)?;
706
707 Ok(())
708}
709
710fn positive(value: Option<i32>) -> Option<i32> {
713 value.filter(|v| *v > 0)
714}
715
716#[cfg(test)]
717mod tests {
718 use super::*;
719
720 use crate::db::connection::Database;
723 use crate::db::queries;
724 use std::sync::{Arc, Mutex};
725
726 fn test_db() -> (Database, tempfile::TempDir) {
727 let dir = tempfile::tempdir().unwrap();
728 let db = Database::open(&dir.path().join("sync_test.db")).unwrap();
729 (db, dir)
730 }
731
732 fn remote_track_meta(remote_id: &str, title: &str, artist: &str, album: &str) -> TrackMeta {
734 TrackMeta {
735 title: title.into(),
736 artist: artist.into(),
737 album_artist: Some(artist.into()),
738 album: album.into(),
739 date: Some("2024".into()),
740 disc: Some(1),
741 track_number: Some(1),
742 genre: Some("Electronic".into()),
743 label: None,
744 duration_ms: Some(240_000),
745 codec: Some("FLAC".into()),
746 sample_rate: Some(44100),
747 bit_depth: Some(16),
748 channels: Some(2),
749 bitrate: Some(1000),
750 size_bytes: None,
751 mtime: None,
752 path: None,
753 source: "remote".into(),
754 remote_id: Some(remote_id.into()),
755 remote_url: Some(format!("https://example.com/stream?id={}", remote_id)),
756 album_remote_id: Some(format!("album-of-{remote_id}")),
757 artist_remote_id: Some(format!("artist-of-{remote_id}")),
758 mbid: Some(format!("mbid-of-{remote_id}")),
759 album_mbid: None,
760 album_added_at: None,
761 }
762 }
763
764 #[test]
767 fn a_zero_quality_figure_is_treated_as_absent() {
768 assert_eq!(positive(Some(0)), None);
769 assert_eq!(positive(Some(16)), Some(16));
770 assert_eq!(positive(None), None);
771 }
772
773 #[test]
777 fn album_metadata_from_the_server_is_recorded() {
778 let (db, _dir) = test_db();
779
780 let meta = remote_track_meta("remote-300", "Anguish", "Sleep", "Volume One");
781 queries::upsert_track(&db.conn, &meta).unwrap();
782
783 queries::enrich_remote_album(
784 &db.conn,
785 "album-of-remote-300",
786 Some("mb-album-1"),
787 Some("volume one"),
788 Some(6),
789 Some("Off The Disk"),
790 )
791 .unwrap();
792
793 let row: (Option<String>, Option<String>, Option<i32>, Option<String>) = db
794 .conn
795 .query_row(
796 "SELECT mbid, sort_name, total_tracks, label FROM albums WHERE title = 'Volume One'",
797 [],
798 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
799 )
800 .unwrap();
801 assert_eq!(row.0.as_deref(), Some("mb-album-1"));
802 assert_eq!(row.1.as_deref(), Some("volume one"));
803 assert_eq!(row.2, Some(6));
804 assert_eq!(row.3.as_deref(), Some("Off The Disk"));
805 }
806
807 #[test]
810 fn server_metadata_does_not_overwrite_what_tags_said() {
811 let (db, _dir) = test_db();
812
813 let mut meta = remote_track_meta("remote-301", "Dopesmoker", "Sleep", "Dopesmoker");
814 meta.label = Some("From The Tags".into());
815 queries::upsert_track(&db.conn, &meta).unwrap();
816
817 queries::enrich_remote_album(
818 &db.conn,
819 "album-of-remote-301",
820 None,
821 None,
822 None,
823 Some("From The Server"),
824 )
825 .unwrap();
826
827 let label: Option<String> = db
828 .conn
829 .query_row(
830 "SELECT label FROM albums WHERE title = 'Dopesmoker'",
831 [],
832 |r| r.get(0),
833 )
834 .unwrap();
835 assert_eq!(label.as_deref(), Some("From The Tags"));
836 }
837
838 #[test]
841 fn artist_metadata_from_the_server_is_recorded() {
842 let (db, _dir) = test_db();
843
844 let meta = remote_track_meta("remote-302", "Holy Mountain", "Sleep", "Holy Mountain");
845 queries::upsert_track(&db.conn, &meta).unwrap();
846
847 queries::enrich_remote_artist(
848 &db.conn,
849 "artist-of-remote-302",
850 Some("mb-artist-1"),
851 Some("sleep"),
852 )
853 .unwrap();
854
855 let row: (Option<String>, Option<String>) = db
856 .conn
857 .query_row(
858 "SELECT mbid, sort_name FROM artists WHERE name = 'Sleep'",
859 [],
860 |r| Ok((r.get(0)?, r.get(1)?)),
861 )
862 .unwrap();
863 assert_eq!(row.0.as_deref(), Some("mb-artist-1"));
864 assert_eq!(row.1.as_deref(), Some("sleep"));
865 }
866
867 #[test]
869 fn a_synced_track_keeps_its_musicbrainz_id() {
870 let (db, _dir) = test_db();
871
872 let meta = remote_track_meta("remote-303", "Aquarian", "Sleep", "Dopesmoker");
873 let id = queries::upsert_track(&db.conn, &meta).unwrap();
874
875 let mbid: Option<String> = db
876 .conn
877 .query_row("SELECT mbid FROM tracks WHERE id = ?1", [id], |r| r.get(0))
878 .unwrap();
879 assert_eq!(mbid.as_deref(), Some("mbid-of-remote-303"));
880 }
881
882 #[test]
886 fn a_synced_track_keeps_its_quality_figures() {
887 let (db, _dir) = test_db();
888
889 let meta = remote_track_meta("remote-200", "Anguish", "Sleep", "Volume One");
890 let id = queries::upsert_track(&db.conn, &meta).unwrap();
891
892 let row = queries::get_track_row(&db.conn, id).unwrap().unwrap();
893 assert_eq!(row.sample_rate, Some(44100));
894 assert_eq!(row.bit_depth, Some(16));
895 assert_eq!(row.channels, Some(2));
896 }
897
898 #[test]
902 fn a_sync_records_the_album_and_artist_ids_too() {
903 let (db, _dir) = test_db();
904
905 let meta = remote_track_meta("remote-100", "Enter", "Russian Circles", "Enter");
906 queries::upsert_track(&db.conn, &meta).unwrap();
907
908 let album: Option<String> = db
909 .conn
910 .query_row(
911 "SELECT remote_id FROM albums WHERE title = 'Enter'",
912 [],
913 |r| r.get(0),
914 )
915 .unwrap();
916 assert_eq!(album.as_deref(), Some("album-of-remote-100"));
917
918 let artist: Option<String> = db
919 .conn
920 .query_row(
921 "SELECT remote_id FROM artists WHERE name = 'Russian Circles'",
922 [],
923 |r| r.get(0),
924 )
925 .unwrap();
926 assert_eq!(artist.as_deref(), Some("artist-of-remote-100"));
927 }
928
929 #[test]
930 fn sync_upserts_tracks_to_database() {
931 let (db, _dir) = test_db();
932
933 let meta = remote_track_meta("remote-001", "Vordhosbn", "Aphex Twin", "Drukqs");
934 let track_id = queries::upsert_track(&db.conn, &meta).unwrap();
935 assert!(track_id > 0, "upsert should return a valid track ID");
936
937 let row = queries::get_track_row(&db.conn, track_id)
939 .unwrap()
940 .expect("track should exist in DB");
941 assert_eq!(row.title, "Vordhosbn");
942 assert_eq!(row.artist_name, "Aphex Twin");
943 assert_eq!(row.album_title, "Drukqs");
944 assert_eq!(row.remote_id.as_deref(), Some("remote-001"));
945 assert_eq!(row.source, "remote");
946 }
947
948 struct StubServer {
953 addr: std::net::SocketAddr,
954 shutdown: Arc<std::sync::atomic::AtomicBool>,
955 }
956
957 #[derive(Default)]
958 struct StubState {
959 albums: Mutex<Vec<(String, String, String)>>,
961 failing: Mutex<HashSet<String>>,
963 unlisted: Mutex<HashSet<String>>,
966 insert_after_first_page: Mutex<Option<(String, String, String)>>,
969 list_pages_served: AtomicUsize,
970 list_types: Mutex<Vec<String>>,
971 album_calls: Mutex<Vec<String>>,
972 lists_songs: bool,
976 song_pages: Mutex<Vec<usize>>,
978 failing_song_pages: Mutex<HashSet<usize>>,
980 library_modified: Mutex<Option<i64>>,
982 }
983
984 impl StubServer {
985 fn start(state: Arc<StubState>) -> Self {
986 let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
987 listener.set_nonblocking(true).unwrap();
988 let addr = listener.local_addr().unwrap();
989 let shutdown = Arc::new(std::sync::atomic::AtomicBool::new(false));
990
991 let stop = shutdown.clone();
992 std::thread::spawn(move || {
993 while !stop.load(Ordering::Relaxed) {
994 match listener.accept() {
995 Ok((stream, _)) => {
996 let _ = stream.set_nonblocking(false);
998 let state = state.clone();
999 std::thread::spawn(move || handle(stream, state));
1000 }
1001 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
1002 std::thread::sleep(std::time::Duration::from_millis(2));
1003 }
1004 Err(_) => break,
1005 }
1006 }
1007 });
1008
1009 Self { addr, shutdown }
1010 }
1011
1012 fn url(&self) -> String {
1013 format!("http://{}", self.addr)
1014 }
1015 }
1016
1017 impl Drop for StubServer {
1018 fn drop(&mut self) {
1019 self.shutdown.store(true, Ordering::Relaxed);
1020 }
1021 }
1022
1023 fn handle(mut stream: std::net::TcpStream, state: Arc<StubState>) {
1027 use std::io::{BufRead, Write};
1028
1029 let Ok(peek) = stream.try_clone() else { return };
1030 let mut reader = std::io::BufReader::new(peek);
1031
1032 loop {
1033 let mut request_line = String::new();
1034 if reader.read_line(&mut request_line).unwrap_or(0) == 0 {
1035 return;
1036 }
1037 let mut line = String::new();
1038 while reader.read_line(&mut line).unwrap_or(0) > 0 {
1039 if line == "\r\n" || line == "\n" {
1040 break;
1041 }
1042 line.clear();
1043 }
1044
1045 let target = request_line.split_whitespace().nth(1).unwrap_or("/");
1046 let (path, query) = target.split_once('?').unwrap_or((target, ""));
1047 let params: std::collections::HashMap<&str, &str> = query
1048 .split('&')
1049 .filter_map(|kv| kv.split_once('='))
1050 .collect();
1051
1052 let (status, body) = respond(&state, path, ¶ms);
1053
1054 let write = write!(
1055 stream,
1056 "HTTP/1.1 {} OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n",
1057 status,
1058 body.len()
1059 )
1060 .and_then(|()| stream.write_all(body.as_bytes()))
1061 .and_then(|()| stream.flush());
1062 if write.is_err() {
1063 return;
1064 }
1065 }
1066 }
1067
1068 fn respond(
1069 state: &Arc<StubState>,
1070 path: &str,
1071 params: &std::collections::HashMap<&str, &str>,
1072 ) -> (u16, String) {
1073 match path.rsplit('/').next().unwrap_or("") {
1074 "getArtists" => (
1075 200,
1076 r#"{"subsonic-response":{"status":"ok","artists":{"index":[{"artist":[{"id":"ar1","name":"Stub Artist"}]}]}}}"#.to_string(),
1077 ),
1078 "getAlbumList2" => {
1079 state
1080 .list_types
1081 .lock()
1082 .unwrap()
1083 .push(params.get("type").copied().unwrap_or("").to_string());
1084 let offset: usize = params.get("offset").and_then(|o| o.parse().ok()).unwrap_or(0);
1085 let size: usize = params.get("size").and_then(|s| s.parse().ok()).unwrap_or(500);
1086
1087 let albums = state.albums.lock().unwrap();
1088 let unlisted = state.unlisted.lock().unwrap();
1089 let slice: Vec<String> = albums
1090 .iter()
1091 .filter(|(id, _, _)| !unlisted.contains(id))
1092 .skip(offset)
1093 .take(size)
1094 .map(|(id, name, created)| {
1095 format!(
1096 r#"{{"id":"{}","name":"{}","artist":"Stub Artist","created":"{}","songCount":1}}"#,
1097 id, name, created
1098 )
1099 })
1100 .collect();
1101 drop(albums);
1102 drop(unlisted);
1103
1104 if state.list_pages_served.fetch_add(1, Ordering::SeqCst) == 0
1105 && let Some(new_album) = state.insert_after_first_page.lock().unwrap().take()
1106 {
1107 state.albums.lock().unwrap().insert(0, new_album);
1108 }
1109
1110 (
1111 200,
1112 format!(
1113 r#"{{"subsonic-response":{{"status":"ok","albumList2":{{"album":[{}]}}}}}}"#,
1114 slice.join(",")
1115 ),
1116 )
1117 }
1118 "getIndexes" => match *state.library_modified.lock().unwrap() {
1119 Some(at) => (
1120 200,
1121 format!(
1122 r#"{{"subsonic-response":{{"status":"ok","indexes":{{"lastModified":{at}}}}}}}"#
1123 ),
1124 ),
1125 None => (200, r#"{"subsonic-response":{"status":"ok"}}"#.to_string()),
1126 },
1127 "search3" => {
1128 let offset: usize = params.get("songOffset").and_then(|o| o.parse().ok()).unwrap_or(0);
1129 let size: usize = params.get("songCount").and_then(|s| s.parse().ok()).unwrap_or(20);
1130 state.song_pages.lock().unwrap().push(offset);
1131 if state.failing_song_pages.lock().unwrap().contains(&offset) {
1132 return (500, r#"{"error":"boom"}"#.to_string());
1133 }
1134 let songs: Vec<String> = if state.lists_songs {
1135 state
1136 .albums
1137 .lock()
1138 .unwrap()
1139 .iter()
1140 .skip(offset)
1141 .take(size)
1142 .map(|(id, name, _)| {
1143 format!(
1144 r#"{{"id":"s{id}","title":"Song {id}","albumId":"{id}","album":"{name}","track":1,"suffix":"flac"}}"#
1145 )
1146 })
1147 .collect()
1148 } else {
1149 Vec::new()
1150 };
1151 (
1152 200,
1153 format!(
1154 r#"{{"subsonic-response":{{"status":"ok","searchResult3":{{"song":[{}]}}}}}}"#,
1155 songs.join(",")
1156 ),
1157 )
1158 }
1159 "getAlbum" => {
1160 let id = params.get("id").copied().unwrap_or("");
1161 state.album_calls.lock().unwrap().push(id.to_string());
1162 if state.failing.lock().unwrap().contains(id) {
1163 (500, r#"{"error":"boom"}"#.to_string())
1164 } else if !state.albums.lock().unwrap().iter().any(|(a, _, _)| a == id) {
1165 (
1166 200,
1167 r#"{"subsonic-response":{"status":"failed","error":{"code":70,"message":"not found"}}}"#
1168 .to_string(),
1169 )
1170 } else {
1171 (
1172 200,
1173 format!(
1174 r#"{{"subsonic-response":{{"status":"ok","album":{{"id":"{id}","name":"{name}","artist":"Stub Artist","song":[{{"id":"s{id}","title":"Song {id}","track":1,"suffix":"flac"}}]}}}}}}"#,
1175 name = state
1176 .albums
1177 .lock()
1178 .unwrap()
1179 .iter()
1180 .find(|(a, _, _)| a == id)
1181 .map_or_else(|| format!("Album {id}"), |(_, n, _)| n.clone())
1182 ),
1183 )
1184 }
1185 }
1186 _ => (404, "{}".to_string()),
1187 }
1188 }
1189
1190 fn stub_albums(n: usize) -> Vec<(String, String, String)> {
1191 (0..n)
1192 .map(|i| {
1193 (
1194 format!("a{:04}", i),
1195 format!("Album {:04}", i),
1196 "2024-01-15T10:30:00Z".to_string(),
1197 )
1198 })
1199 .collect()
1200 }
1201
1202 #[test]
1203 fn failed_album_fetch_does_not_advance_last_sync_and_next_sync_retries() {
1204 let (db, _dir) = test_db();
1205 let state = Arc::new(StubState {
1206 albums: Mutex::new(stub_albums(4)),
1207 failing: Mutex::new(["a0002".to_string()].into_iter().collect()),
1208 ..Default::default()
1209 });
1210 let server = StubServer::start(state.clone());
1211 let client = SubsonicClient::new(&server.url(), "u", "p");
1212
1213 let first = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1214 assert_eq!(first.albums_failed, 1, "the failing album must be counted");
1215 assert_eq!(first.albums_synced, 3);
1216 assert!(!first.is_complete());
1217 assert_eq!(
1218 get_last_sync(&db, &server.url()).unwrap(),
1219 None,
1220 "an incomplete sync must not advance last_sync"
1221 );
1222
1223 state.failing.lock().unwrap().clear();
1226 state.album_calls.lock().unwrap().clear();
1227 let second = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1228
1229 assert!(
1230 state
1231 .album_calls
1232 .lock()
1233 .unwrap()
1234 .contains(&"a0002".to_string()),
1235 "the previously failed album must be retried"
1236 );
1237 assert_eq!(second.albums_failed, 0);
1238 assert!(second.is_complete());
1239 assert!(
1240 get_last_sync(&db, &server.url()).unwrap().is_some(),
1241 "a clean sync advances last_sync"
1242 );
1243 }
1244
1245 #[test]
1246 fn album_inserted_mid_pagination_is_not_fetched_twice_or_skipped() {
1247 let (db, _dir) = test_db();
1251 let state = Arc::new(StubState {
1252 albums: Mutex::new(stub_albums(600)),
1253 insert_after_first_page: Mutex::new(Some((
1254 "aNEW".to_string(),
1255 "AAA Brand New".to_string(),
1256 "2024-06-01T00:00:00Z".to_string(),
1257 ))),
1258 ..Default::default()
1259 });
1260 let server = StubServer::start(state.clone());
1261 let client = SubsonicClient::new(&server.url(), "u", "p");
1262
1263 let result = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1264 assert_eq!(result.albums_failed, 0);
1265
1266 let calls = state.album_calls.lock().unwrap().clone();
1267 let unique: HashSet<&String> = calls.iter().collect();
1268 assert_eq!(
1269 calls.len(),
1270 unique.len(),
1271 "no album may be fetched twice after the window shifts"
1272 );
1273
1274 for i in 0..600 {
1276 let id = format!("a{:04}", i);
1277 assert!(unique.contains(&id), "album {} was skipped", id);
1278 }
1279
1280 let types = state.list_types.lock().unwrap().clone();
1281 assert!(
1282 types.iter().all(|t| t == "alphabeticalByName"),
1283 "the paginated walk must use a stable ordering, got {:?}",
1284 types
1285 );
1286 }
1287
1288 #[test]
1289 fn a_full_sync_relinks_files_whose_server_id_changed() {
1290 let (db, _dir) = test_db();
1291 let file = TrackMeta {
1292 date: None,
1293 disc: None,
1294 path: Some("/music/song.flac".into()),
1295 source: "local".into(),
1296 album_remote_id: None,
1297 artist_remote_id: None,
1298 mbid: None,
1299 ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1300 };
1301 queries::upsert_track(&db.conn, &file).unwrap();
1302
1303 let state = Arc::new(StubState {
1304 albums: Mutex::new(stub_albums(1)),
1305 ..Default::default()
1306 });
1307 let server = StubServer::start(state);
1308 let client = SubsonicClient::new(&server.url(), "u", "p");
1309 sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1310
1311 let rows: Vec<(Option<String>, Option<String>)> = db
1312 .conn
1313 .prepare("SELECT path, remote_id FROM tracks WHERE title = 'Song a0000'")
1314 .unwrap()
1315 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1316 .unwrap()
1317 .collect::<Result<_, _>>()
1318 .unwrap();
1319 assert_eq!(
1320 rows,
1321 vec![(Some("/music/song.flac".into()), Some("sa0000".into()))],
1322 "the file takes the server's current id, and the recording is listed once"
1323 );
1324 }
1325
1326 #[test]
1327 fn a_full_sync_folds_server_rows_whose_id_vanished() {
1328 let (db, _dir) = test_db();
1329 let ghost = TrackMeta {
1330 date: None,
1331 disc: None,
1332 album_remote_id: None,
1333 artist_remote_id: None,
1334 mbid: None,
1335 ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1336 };
1337 let ghost_id = queries::upsert_track(&db.conn, &ghost).unwrap();
1338 let ghost_url = ghost.remote_url.clone().unwrap();
1339 db.conn
1340 .execute(
1341 "INSERT INTO favourites (track_path) VALUES (?1)",
1342 params![ghost_url],
1343 )
1344 .unwrap();
1345
1346 let state = Arc::new(StubState {
1347 albums: Mutex::new(stub_albums(1)),
1348 ..Default::default()
1349 });
1350 let server = StubServer::start(state);
1351 let client = SubsonicClient::new(&server.url(), "u", "p");
1352 sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1353
1354 let rows: Vec<(i64, String, String)> = db
1355 .conn
1356 .prepare("SELECT id, remote_id, remote_url FROM tracks WHERE title = 'Song a0000'")
1357 .unwrap()
1358 .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))
1359 .unwrap()
1360 .collect::<Result<_, _>>()
1361 .unwrap();
1362 assert_eq!(rows.len(), 1, "the dead row takes the live id");
1363 let (id, remote_id, remote_url) = &rows[0];
1364 assert_eq!(*id, ghost_id);
1365 assert_eq!(remote_id, "sa0000");
1366 let favourites: Vec<String> = db
1367 .conn
1368 .prepare("SELECT track_path FROM favourites")
1369 .unwrap()
1370 .query_map([], |row| row.get(0))
1371 .unwrap()
1372 .collect::<Result<_, _>>()
1373 .unwrap();
1374 assert_eq!(
1375 favourites,
1376 vec![remote_url.clone()],
1377 "the favourite follows"
1378 );
1379 }
1380
1381 #[test]
1384 fn a_library_that_has_not_moved_is_not_walked_again() {
1385 use crate::helpers::{Walk, sync_remote};
1386 let (db, _dir) = test_db();
1387 let state = Arc::new(StubState {
1388 albums: Mutex::new(stub_albums(3)),
1389 library_modified: Mutex::new(Some(1_000)),
1390 ..Default::default()
1391 });
1392 let server = StubServer::start(state.clone());
1393 let client = SubsonicClient::new(&server.url(), "u", "p");
1394 let walks = || state.list_types.lock().unwrap().len();
1395 let sync = |walk| sync_remote(&db, &client, walk, &server.url(), "u", &|_| {}).unwrap();
1396
1397 sync(Walk::IfChanged);
1398 let first = walks();
1399 assert!(first > 0, "nothing recorded, so walked");
1400
1401 sync(Walk::IfChanged);
1402 assert_eq!(walks(), first, "unchanged, so not walked");
1403
1404 sync(Walk::Always);
1405 let asked = walks();
1406 assert!(asked > first, "asked for, so walked");
1407
1408 *state.library_modified.lock().unwrap() = Some(2_000);
1409 sync(Walk::IfChanged);
1410 assert!(walks() > asked, "moved, so walked");
1411 }
1412
1413 #[test]
1415 fn a_server_without_a_library_version_is_always_walked() {
1416 use crate::helpers::{Walk, sync_remote};
1417 let (db, _dir) = test_db();
1418 let state = Arc::new(StubState {
1419 albums: Mutex::new(stub_albums(3)),
1420 ..Default::default()
1421 });
1422 let server = StubServer::start(state.clone());
1423 let client = SubsonicClient::new(&server.url(), "u", "p");
1424 let walks = || state.list_types.lock().unwrap().len();
1425
1426 sync_remote(&db, &client, Walk::IfChanged, &server.url(), "u", &|_| {}).unwrap();
1427 let first = walks();
1428 sync_remote(&db, &client, Walk::IfChanged, &server.url(), "u", &|_| {}).unwrap();
1429 assert!(walks() > first);
1430 }
1431
1432 #[test]
1435 fn an_album_left_out_of_the_listing_is_removed_only_when_gone() {
1436 let (db, _dir) = test_db();
1437 let state = Arc::new(StubState {
1438 albums: Mutex::new(stub_albums(3)),
1439 ..Default::default()
1440 });
1441 let server = StubServer::start(state.clone());
1442 let client = SubsonicClient::new(&server.url(), "u", "p");
1443 sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1444
1445 state.unlisted.lock().unwrap().insert("a0001".into());
1446 let shifted = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1447 assert_eq!(shifted.tracks_removed, 0, "still there, only unlisted");
1448
1449 state
1450 .albums
1451 .lock()
1452 .unwrap()
1453 .retain(|(id, _, _)| id != "a0001");
1454 let deleted = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1455 assert_eq!(deleted.tracks_removed, 1);
1456 }
1457
1458 #[test]
1461 fn a_full_sync_pages_songs_instead_of_fetching_each_album() {
1462 let (db, _dir) = test_db();
1463 let state = Arc::new(StubState {
1464 albums: Mutex::new(stub_albums(1_234)),
1465 lists_songs: true,
1466 ..Default::default()
1467 });
1468 let server = StubServer::start(state.clone());
1469 let client = SubsonicClient::new(&server.url(), "u", "p");
1470
1471 let seen = Mutex::new(Vec::new());
1472 let result = sync_library(&db, &client, &server.url(), "u", &|p| {
1473 seen.lock().unwrap().push(p)
1474 })
1475 .unwrap();
1476
1477 assert!(state.album_calls.lock().unwrap().is_empty(), "no getAlbum");
1478 let mut pages = state.song_pages.lock().unwrap().clone();
1479 pages.sort();
1480 pages.dedup();
1481 assert_eq!(pages[..3], [0, 500, 1000]);
1482 assert_eq!(result.tracks_synced, 1_234);
1483 assert_eq!(result.albums_synced, 1_234);
1484 assert!(result.is_complete());
1485 assert!(get_last_sync(&db, &server.url()).unwrap().is_some());
1486
1487 let album: (Option<String>, Option<i32>) = db
1488 .conn
1489 .query_row(
1490 "SELECT al.remote_id, al.total_tracks FROM albums al WHERE al.title = 'Album 0007'",
1491 [],
1492 |r| Ok((r.get(0)?, r.get(1)?)),
1493 )
1494 .unwrap();
1495 assert_eq!(
1496 album,
1497 (Some("a0007".into()), Some(1)),
1498 "album metadata comes from the list"
1499 );
1500
1501 let seen = seen.into_inner().unwrap();
1502 let tracks: Vec<&SyncProgress> = seen
1503 .iter()
1504 .filter(|p| p.phase == SyncPhase::Tracks)
1505 .collect();
1506 assert!(tracks.len() >= 3, "progress at least per page");
1507 assert!(tracks.iter().all(|p| p.total == Some(1_234)));
1508 assert_eq!(tracks.last().unwrap().done, 1_234);
1509 assert_eq!(seen.last().unwrap().phase, SyncPhase::Finishing);
1510 }
1511
1512 #[test]
1515 fn a_server_that_lists_no_songs_is_synced_album_by_album() {
1516 let (db, _dir) = test_db();
1517 let state = Arc::new(StubState {
1518 albums: Mutex::new(stub_albums(3)),
1519 ..Default::default()
1520 });
1521 let server = StubServer::start(state.clone());
1522 let client = SubsonicClient::new(&server.url(), "u", "p");
1523
1524 let result = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1525
1526 assert_eq!(state.song_pages.lock().unwrap().clone(), vec![0]);
1527 assert_eq!(state.album_calls.lock().unwrap().len(), 3);
1528 assert_eq!(result.tracks_synced, 3);
1529 assert!(result.is_complete());
1530 }
1531
1532 #[test]
1535 fn a_failed_song_page_does_not_advance_last_sync() {
1536 let (db, _dir) = test_db();
1537 let state = Arc::new(StubState {
1538 albums: Mutex::new(stub_albums(1_200)),
1539 lists_songs: true,
1540 failing_song_pages: Mutex::new([500].into_iter().collect()),
1541 ..Default::default()
1542 });
1543 let server = StubServer::start(state.clone());
1544 let client = SubsonicClient::new(&server.url(), "u", "p");
1545
1546 let result = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1547
1548 assert_eq!(result.pages_failed, 1);
1549 assert!(!result.is_complete());
1550 assert_eq!(get_last_sync(&db, &server.url()).unwrap(), None);
1551 let retries = state
1552 .song_pages
1553 .lock()
1554 .unwrap()
1555 .iter()
1556 .filter(|&&o| o == 500)
1557 .count();
1558 assert_eq!(retries as u32, PAGE_ATTEMPTS);
1559 }
1560
1561 #[test]
1562 fn sync_deduplicates_by_remote_id() {
1563 let (db, _dir) = test_db();
1564
1565 let meta1 = remote_track_meta("remote-dup", "Original Title", "Artist A", "Album X");
1567 let id1 = queries::upsert_track(&db.conn, &meta1).unwrap();
1568
1569 let meta2 = remote_track_meta("remote-dup", "Updated Title", "Artist A", "Album X");
1571 let id2 = queries::upsert_track(&db.conn, &meta2).unwrap();
1572
1573 assert_eq!(id1, id2, "same remote_id should resolve to same track row");
1575
1576 let row = queries::get_track_row(&db.conn, id2)
1578 .unwrap()
1579 .expect("track should exist");
1580 assert_eq!(row.title, "Updated Title");
1581 assert_eq!(row.remote_id.as_deref(), Some("remote-dup"));
1582
1583 let stats = queries::library_stats(&db.conn).unwrap();
1585 assert_eq!(
1586 stats.total_tracks, 1,
1587 "should have exactly 1 track after dedup"
1588 );
1589 }
1590
1591 #[test]
1595 fn resyncing_a_track_under_a_new_server_id_keeps_one_row() {
1596 let (db, _dir) = test_db();
1597 let before = remote_track_meta("46215", "Archangel", "Burial", "Untrue");
1598 let row = queries::upsert_synced_track(&db.conn, &before, &HashSet::from(["46215".into()]))
1599 .unwrap();
1600 queries::add_favourite(
1601 &db.conn,
1602 queries::LOCAL_USER,
1603 std::path::Path::new(before.remote_url.as_deref().unwrap()),
1604 )
1605 .unwrap();
1606
1607 let uid = "0199a0b2-7c4e-7d3a-9f1b-2c3d4e5f6a7b";
1608 let after = remote_track_meta(uid, "Archangel", "Burial", "Untrue");
1609 let again =
1610 queries::upsert_synced_track(&db.conn, &after, &HashSet::from([uid.into()])).unwrap();
1611
1612 assert_eq!(again, row, "the same row, not a second copy");
1613 assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 1);
1614 let (remote_id, album, artist): (String, String, String) = db
1615 .conn
1616 .query_row(
1617 "SELECT t.remote_id, al.remote_id, ar.remote_id FROM tracks t
1618 JOIN albums al ON al.id = t.album_id JOIN artists ar ON ar.id = t.artist_id",
1619 [],
1620 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1621 )
1622 .unwrap();
1623 assert_eq!(remote_id, uid);
1624 assert_eq!(album, format!("album-of-{uid}"));
1625 assert_eq!(artist, format!("artist-of-{uid}"));
1626 let adopted: String = db
1627 .conn
1628 .query_row("SELECT uid FROM tracks WHERE id = ?1", [row], |r| r.get(0))
1629 .unwrap();
1630 assert_eq!(adopted, uid, "the row takes the server's uid");
1631 let favourites = queries::load_favourites(&db.conn, queries::LOCAL_USER).unwrap();
1632 assert_eq!(
1633 favourites,
1634 HashSet::from([std::path::PathBuf::from(after.remote_url.unwrap())]),
1635 "the favourite follows the new stream address"
1636 );
1637 }
1638
1639 #[test]
1642 fn identical_entries_in_one_sync_stay_two_rows() {
1643 let (db, _dir) = test_db();
1644 let mut seen = HashSet::new();
1645 for id in ["dup-1", "dup-2"] {
1646 seen.insert(id.to_string());
1647 queries::upsert_synced_track(
1648 &db.conn,
1649 &remote_track_meta(id, "Archangel", "Burial", "Untrue"),
1650 &seen,
1651 )
1652 .unwrap();
1653 }
1654 assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 2);
1655 }
1656}