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);
250 let live_tracks =
251 (result.is_complete() && !song_ids.is_empty() && listed_everything).then_some(&song_ids);
252 if result.is_complete() && !live_albums.is_empty() {
253 match queries::remove_vanished_remote(&db.conn, live_tracks, Some(&live_albums)) {
254 Ok(0) => {}
255 Ok(n) => {
256 result.tracks_removed = n;
257 log::info!("{n} tracks the server no longer has; removed");
258 }
259 Err(e) => log::warn!("failed to remove tracks the server deleted: {e}"),
260 }
261 }
262
263 if result.is_complete() {
264 update_last_sync(db, server_url, username, sync_start)?;
265 } else {
266 log::warn!(
267 "{} album(s) and {} page(s) failed to fetch and {} track(s) failed to write — the next sync walks the library again",
268 result.albums_failed,
269 result.pages_failed,
270 result.tracks_failed,
271 );
272 }
273
274 log::info!(
275 "sync complete: {} artists, {} albums, {} tracks, {} albums and {} tracks failed",
276 result.artists_synced,
277 result.albums_synced,
278 result.tracks_synced,
279 result.albums_failed,
280 result.tracks_failed,
281 );
282
283 db.optimize();
284
285 Ok(result)
286}
287
288fn confirm_missing_albums(db: &Database, client: &SubsonicClient, live: &mut HashSet<String>) {
294 let held: Vec<String> = db
295 .conn
296 .prepare(
297 "SELECT DISTINCT al.remote_id FROM albums al JOIN tracks t ON t.album_id = al.id
298 WHERE al.remote_id IS NOT NULL AND t.path IS NULL AND t.remote_id IS NOT NULL",
299 )
300 .and_then(|mut stmt| {
301 stmt.query_map([], |r| r.get(0))?
302 .collect::<rusqlite::Result<Vec<String>>>()
303 })
304 .unwrap_or_default();
305 let missing: Vec<String> = held.into_iter().filter(|id| !live.contains(id)).collect();
306 for id in missing {
307 match client.get_album(&id) {
308 Err(SubsonicError::Api { code: 70, .. }) => {}
309 Ok(_) => {
310 log::info!("album {id} was missing from the listing but still exists");
311 live.insert(id);
312 }
313 Err(e) => {
314 log::warn!("could not confirm album {id} was deleted ({e}); keeping it");
315 live.insert(id);
316 }
317 }
318 }
319}
320
321fn list_albums(
323 client: &SubsonicClient,
324 progress: &(dyn Fn(SyncProgress) + Sync),
325) -> Result<Vec<SubsonicAlbum>, SyncError> {
326 let mut albums = Vec::new();
327 let mut seen: HashSet<String> = HashSet::new();
330 let mut offset = 0u32;
331 loop {
332 let page = client.get_album_list("alphabeticalByName", PAGE_SIZE, offset)?;
333 let count = page.len() as u32;
334 offset += count;
335 albums.extend(page.into_iter().filter(|a| seen.insert(a.id.clone())));
336 progress(SyncProgress {
337 phase: SyncPhase::Albums,
338 done: albums.len() as u64,
339 total: None,
340 });
341 if count < PAGE_SIZE {
342 return Ok(albums);
343 }
344 }
345}
346
347fn sync_all_songs(
354 db: &Database,
355 client: &SubsonicClient,
356 albums: &[SubsonicAlbum],
357 total: Option<u64>,
358 result: &mut SyncResult,
359 song_ids: &mut HashSet<String>,
360 progress: &(dyn Fn(SyncProgress) + Sync),
361) -> Result<bool, SyncError> {
362 let first = match client.all_songs_page(PAGE_SIZE, 0) {
363 Ok(songs) if !songs.is_empty() => songs,
364 Ok(_) => {
365 log::info!("server lists no songs for an empty search; syncing album by album");
366 return Ok(false);
367 }
368 Err(e) => {
369 log::info!("server refused an empty search ({e}); syncing album by album");
370 return Ok(false);
371 }
372 };
373
374 let by_id: HashMap<&str, &SubsonicAlbum> = albums.iter().map(|a| (a.id.as_str(), a)).collect();
375 let mut albums_seen: HashSet<String> = HashSet::new();
376 let mut write = |songs: Vec<SubsonicSong>,
377 result: &mut SyncResult,
378 song_ids: &mut HashSet<String>|
379 -> Result<(), SyncError> {
380 let batch = group_by_album(songs, &by_id);
381 albums_seen.extend(batch.iter().map(|a| a.id.clone()));
382 write_albums(db, client, &batch, result, song_ids)?;
383 progress(SyncProgress {
384 phase: SyncPhase::Tracks,
385 done: song_ids.len() as u64,
386 total,
387 });
388 Ok(())
389 };
390
391 let short = (first.len() as u32) < PAGE_SIZE;
392 write(first, result, song_ids)?;
393
394 if !short {
395 let next = AtomicU32::new(PAGE_SIZE);
401 let end = AtomicU32::new(u32::MAX);
404 let (tx, rx) = std::sync::mpsc::sync_channel(FETCH_LANES);
405 let failed = std::thread::scope(|scope| -> Result<usize, SyncError> {
406 for _ in 0..FETCH_LANES {
407 let tx = tx.clone();
408 let (next, end) = (&next, &end);
409 scope.spawn(move || {
410 loop {
411 let offset = next.fetch_add(PAGE_SIZE, Ordering::Relaxed);
412 if offset >= end.load(Ordering::Relaxed) {
413 return;
414 }
415 let page = fetch_page(client, offset);
416 if page
421 .as_ref()
422 .map_or(true, |songs| (songs.len() as u32) < PAGE_SIZE)
423 {
424 end.fetch_min(offset, Ordering::Relaxed);
425 }
426 if tx.send((offset, page)).is_err() {
427 return;
428 }
429 }
430 });
431 }
432 drop(tx);
433
434 let mut failed = 0;
435 for (offset, page) in rx {
436 match page {
437 Ok(songs) => write(songs, result, song_ids)?,
438 Err(e) => {
439 log::warn!("failed to fetch songs from offset {offset}: {e}");
440 failed += 1;
441 }
442 }
443 }
444 Ok(failed)
445 })?;
446 result.pages_failed += failed;
447 }
448
449 result.albums_synced += albums_seen.len();
450 Ok(true)
451}
452
453fn fetch_page(
456 client: &SubsonicClient,
457 offset: u32,
458) -> Result<Vec<SubsonicSong>, super::client::SubsonicError> {
459 let mut attempt = 1;
460 loop {
461 match client.all_songs_page(PAGE_SIZE, offset) {
462 Ok(songs) => return Ok(songs),
463 Err(e) if attempt >= PAGE_ATTEMPTS => return Err(e),
464 Err(e) => {
465 log::debug!("songs from offset {offset}, attempt {attempt}: {e}");
466 std::thread::sleep(std::time::Duration::from_millis(250 * attempt as u64));
467 attempt += 1;
468 }
469 }
470 }
471}
472
473fn group_by_album(
479 songs: Vec<SubsonicSong>,
480 albums: &HashMap<&str, &SubsonicAlbum>,
481) -> Vec<SubsonicAlbumFull> {
482 let mut grouped: Vec<SubsonicAlbumFull> = Vec::new();
483 let mut index: HashMap<String, usize> = HashMap::new();
484 for song in songs {
485 let Some(album_id) = song.album_id.clone() else {
486 log::warn!("song {} has no album id; skipped", song.id);
487 continue;
488 };
489 let i = *index.entry(album_id.clone()).or_insert_with(|| {
490 grouped.push(match albums.get(album_id.as_str()) {
491 Some(album) => SubsonicAlbumFull {
492 id: album.id.clone(),
493 name: album.name.clone(),
494 artist: album.artist.clone(),
495 artist_id: album.artist_id.clone(),
496 year: album.year,
497 genre: album.genre.clone(),
498 song_count: album.song_count,
499 created: album.created.clone(),
500 music_brainz_id: album.music_brainz_id.clone(),
501 sort_name: album.sort_name.clone(),
502 record_labels: album.record_labels.clone(),
503 song: Vec::new(),
504 },
505 None => SubsonicAlbumFull {
506 id: album_id,
507 name: song.album.clone().unwrap_or_default(),
508 artist: song.artist.clone(),
509 artist_id: song.artist_id.clone(),
510 year: song.year,
511 genre: song.genre.clone(),
512 song_count: None,
513 created: None,
514 music_brainz_id: None,
515 sort_name: None,
516 record_labels: Vec::new(),
517 song: Vec::new(),
518 },
519 });
520 grouped.len() - 1
521 });
522 grouped[i].song.push(song);
523 }
524 grouped
525}
526
527const ALBUM_BATCH: usize = 100;
530
531fn sync_by_album(
534 db: &Database,
535 client: &SubsonicClient,
536 albums: &[&SubsonicAlbum],
537 total: Option<u64>,
538 result: &mut SyncResult,
539 song_ids: &mut HashSet<String>,
540 progress: &(dyn Fn(SyncProgress) + Sync),
541) -> Result<(), SyncError> {
542 for batch in albums.chunks(ALBUM_BATCH) {
543 let failures = AtomicUsize::new(0);
544 let fetched: Vec<SubsonicAlbumFull> = batch
545 .par_iter()
546 .filter_map(|album| match client.get_album(&album.id) {
547 Ok(full) => Some(full),
548 Err(e) => {
549 log::warn!("failed to fetch album {}: {}", album.id, e);
550 failures.fetch_add(1, Ordering::Relaxed);
551 None
552 }
553 })
554 .collect();
555 result.albums_failed += failures.into_inner();
556 result.albums_synced += fetched.len();
557
558 write_albums(db, client, &fetched, result, song_ids)?;
559 progress(SyncProgress {
560 phase: SyncPhase::Tracks,
561 done: song_ids.len() as u64,
562 total,
563 });
564 log::info!(
565 "synced {} albums ({} tracks) so far...",
566 result.albums_synced,
567 result.tracks_synced
568 );
569 }
570 Ok(())
571}
572
573fn write_artists(db: &Database, artists: &[SubsonicArtist], result: &mut SyncResult) {
578 if db.conn.execute_batch("BEGIN IMMEDIATE").is_err() {
582 return;
583 }
584 let mut enriched = 0;
585 for artist in artists {
586 match queries::enrich_remote_artist(
587 &db.conn,
588 &artist.id,
589 artist.music_brainz_id.as_deref(),
590 artist.sort_name.as_deref(),
591 ) {
592 Ok(()) => enriched += 1,
593 Err(e) => log::warn!(
594 "failed to record artist metadata for {}: {}",
595 artist.name,
596 e
597 ),
598 }
599 }
600 if db.conn.execute_batch("COMMIT").is_err() {
601 let _ = db.conn.execute_batch("ROLLBACK");
602 return;
603 }
604 result.artists_synced = enriched;
605 log::info!("recorded metadata for {enriched} artists");
606}
607
608fn write_albums(
610 db: &Database,
611 client: &SubsonicClient,
612 albums: &[SubsonicAlbumFull],
613 result: &mut SyncResult,
614 song_ids: &mut HashSet<String>,
615) -> Result<(), SyncError> {
616 db.conn
619 .execute_batch("BEGIN IMMEDIATE")
620 .map_err(crate::db::connection::DbError::from)?;
621
622 for album in albums {
623 let artist_name = album.artist.as_deref().unwrap_or("Unknown Artist");
624
625 for song in &album.song {
626 song_ids.insert(song.id.clone());
627 let meta = TrackMeta {
628 title: song.title.clone(),
629 artist: song
630 .artist
631 .clone()
632 .unwrap_or_else(|| artist_name.to_string()),
633 album_artist: album.artist.clone(),
634 album: album.name.clone(),
635 date: album.year.map(|y| y.to_string()),
636 disc: song.disc_number,
637 track_number: song.track,
638 genre: song.genre.clone().or_else(|| album.genre.clone()),
639 label: None,
640 duration_ms: song.duration.map(|d| d * 1000),
641 codec: song.codec.clone().or_else(|| song.suffix.clone()),
642 sample_rate: positive(song.sampling_rate),
649 bit_depth: positive(song.bit_depth),
650 channels: positive(song.channel_count),
651 bitrate: song.bit_rate,
652 size_bytes: None,
653 mtime: None,
654 path: None,
655 source: "remote".to_string(),
656 remote_id: Some(song.id.clone()),
657 remote_url: Some(client.stream_url_template(&song.id)),
658 album_remote_id: Some(album.id.clone()),
659 artist_remote_id: album.artist_id.clone(),
660 mbid: song.music_brainz_id.clone(),
661 album_mbid: album.music_brainz_id.clone(),
662 album_added_at: album.created.clone(),
663 };
664
665 match queries::upsert_synced_track(&db.conn, &meta, song_ids) {
666 Ok(_) => result.tracks_synced += 1,
667 Err(e) => {
668 result.tracks_failed += 1;
669 log::warn!("failed to insert remote track {}: {}", song.title, e);
670 }
671 }
672 }
673
674 if let Err(e) = queries::enrich_remote_album(
678 &db.conn,
679 &album.id,
680 album.music_brainz_id.as_deref(),
681 album.sort_name.as_deref(),
682 album.song_count,
683 album
684 .record_labels
685 .iter()
686 .map(|l| l.name.as_str())
687 .find(|l| !l.is_empty()),
688 ) {
689 log::warn!("failed to record album metadata for {}: {}", album.name, e);
690 }
691 }
692
693 db.conn
694 .execute_batch("COMMIT")
695 .map_err(crate::db::connection::DbError::from)?;
696
697 Ok(())
698}
699
700fn positive(value: Option<i32>) -> Option<i32> {
703 value.filter(|v| *v > 0)
704}
705
706#[cfg(test)]
707mod tests {
708 use super::*;
709
710 use crate::db::connection::Database;
713 use crate::db::queries;
714 use std::sync::{Arc, Mutex};
715
716 fn test_db() -> (Database, tempfile::TempDir) {
717 let dir = tempfile::tempdir().unwrap();
718 let db = Database::open(&dir.path().join("sync_test.db")).unwrap();
719 (db, dir)
720 }
721
722 fn remote_track_meta(remote_id: &str, title: &str, artist: &str, album: &str) -> TrackMeta {
724 TrackMeta {
725 title: title.into(),
726 artist: artist.into(),
727 album_artist: Some(artist.into()),
728 album: album.into(),
729 date: Some("2024".into()),
730 disc: Some(1),
731 track_number: Some(1),
732 genre: Some("Electronic".into()),
733 label: None,
734 duration_ms: Some(240_000),
735 codec: Some("FLAC".into()),
736 sample_rate: Some(44100),
737 bit_depth: Some(16),
738 channels: Some(2),
739 bitrate: Some(1000),
740 size_bytes: None,
741 mtime: None,
742 path: None,
743 source: "remote".into(),
744 remote_id: Some(remote_id.into()),
745 remote_url: Some(format!("https://example.com/stream?id={}", remote_id)),
746 album_remote_id: Some(format!("album-of-{remote_id}")),
747 artist_remote_id: Some(format!("artist-of-{remote_id}")),
748 mbid: Some(format!("mbid-of-{remote_id}")),
749 album_mbid: None,
750 album_added_at: None,
751 }
752 }
753
754 #[test]
757 fn a_zero_quality_figure_is_treated_as_absent() {
758 assert_eq!(positive(Some(0)), None);
759 assert_eq!(positive(Some(16)), Some(16));
760 assert_eq!(positive(None), None);
761 }
762
763 #[test]
767 fn album_metadata_from_the_server_is_recorded() {
768 let (db, _dir) = test_db();
769
770 let meta = remote_track_meta("remote-300", "Anguish", "Sleep", "Volume One");
771 queries::upsert_track(&db.conn, &meta).unwrap();
772
773 queries::enrich_remote_album(
774 &db.conn,
775 "album-of-remote-300",
776 Some("mb-album-1"),
777 Some("volume one"),
778 Some(6),
779 Some("Off The Disk"),
780 )
781 .unwrap();
782
783 let row: (Option<String>, Option<String>, Option<i32>, Option<String>) = db
784 .conn
785 .query_row(
786 "SELECT mbid, sort_name, total_tracks, label FROM albums WHERE title = 'Volume One'",
787 [],
788 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
789 )
790 .unwrap();
791 assert_eq!(row.0.as_deref(), Some("mb-album-1"));
792 assert_eq!(row.1.as_deref(), Some("volume one"));
793 assert_eq!(row.2, Some(6));
794 assert_eq!(row.3.as_deref(), Some("Off The Disk"));
795 }
796
797 #[test]
800 fn server_metadata_does_not_overwrite_what_tags_said() {
801 let (db, _dir) = test_db();
802
803 let mut meta = remote_track_meta("remote-301", "Dopesmoker", "Sleep", "Dopesmoker");
804 meta.label = Some("From The Tags".into());
805 queries::upsert_track(&db.conn, &meta).unwrap();
806
807 queries::enrich_remote_album(
808 &db.conn,
809 "album-of-remote-301",
810 None,
811 None,
812 None,
813 Some("From The Server"),
814 )
815 .unwrap();
816
817 let label: Option<String> = db
818 .conn
819 .query_row(
820 "SELECT label FROM albums WHERE title = 'Dopesmoker'",
821 [],
822 |r| r.get(0),
823 )
824 .unwrap();
825 assert_eq!(label.as_deref(), Some("From The Tags"));
826 }
827
828 #[test]
831 fn artist_metadata_from_the_server_is_recorded() {
832 let (db, _dir) = test_db();
833
834 let meta = remote_track_meta("remote-302", "Holy Mountain", "Sleep", "Holy Mountain");
835 queries::upsert_track(&db.conn, &meta).unwrap();
836
837 queries::enrich_remote_artist(
838 &db.conn,
839 "artist-of-remote-302",
840 Some("mb-artist-1"),
841 Some("sleep"),
842 )
843 .unwrap();
844
845 let row: (Option<String>, Option<String>) = db
846 .conn
847 .query_row(
848 "SELECT mbid, sort_name FROM artists WHERE name = 'Sleep'",
849 [],
850 |r| Ok((r.get(0)?, r.get(1)?)),
851 )
852 .unwrap();
853 assert_eq!(row.0.as_deref(), Some("mb-artist-1"));
854 assert_eq!(row.1.as_deref(), Some("sleep"));
855 }
856
857 #[test]
859 fn a_synced_track_keeps_its_musicbrainz_id() {
860 let (db, _dir) = test_db();
861
862 let meta = remote_track_meta("remote-303", "Aquarian", "Sleep", "Dopesmoker");
863 let id = queries::upsert_track(&db.conn, &meta).unwrap();
864
865 let mbid: Option<String> = db
866 .conn
867 .query_row("SELECT mbid FROM tracks WHERE id = ?1", [id], |r| r.get(0))
868 .unwrap();
869 assert_eq!(mbid.as_deref(), Some("mbid-of-remote-303"));
870 }
871
872 #[test]
876 fn a_synced_track_keeps_its_quality_figures() {
877 let (db, _dir) = test_db();
878
879 let meta = remote_track_meta("remote-200", "Anguish", "Sleep", "Volume One");
880 let id = queries::upsert_track(&db.conn, &meta).unwrap();
881
882 let row = queries::get_track_row(&db.conn, id).unwrap().unwrap();
883 assert_eq!(row.sample_rate, Some(44100));
884 assert_eq!(row.bit_depth, Some(16));
885 assert_eq!(row.channels, Some(2));
886 }
887
888 #[test]
892 fn a_sync_records_the_album_and_artist_ids_too() {
893 let (db, _dir) = test_db();
894
895 let meta = remote_track_meta("remote-100", "Enter", "Russian Circles", "Enter");
896 queries::upsert_track(&db.conn, &meta).unwrap();
897
898 let album: Option<String> = db
899 .conn
900 .query_row(
901 "SELECT remote_id FROM albums WHERE title = 'Enter'",
902 [],
903 |r| r.get(0),
904 )
905 .unwrap();
906 assert_eq!(album.as_deref(), Some("album-of-remote-100"));
907
908 let artist: Option<String> = db
909 .conn
910 .query_row(
911 "SELECT remote_id FROM artists WHERE name = 'Russian Circles'",
912 [],
913 |r| r.get(0),
914 )
915 .unwrap();
916 assert_eq!(artist.as_deref(), Some("artist-of-remote-100"));
917 }
918
919 #[test]
920 fn sync_upserts_tracks_to_database() {
921 let (db, _dir) = test_db();
922
923 let meta = remote_track_meta("remote-001", "Vordhosbn", "Aphex Twin", "Drukqs");
924 let track_id = queries::upsert_track(&db.conn, &meta).unwrap();
925 assert!(track_id > 0, "upsert should return a valid track ID");
926
927 let row = queries::get_track_row(&db.conn, track_id)
929 .unwrap()
930 .expect("track should exist in DB");
931 assert_eq!(row.title, "Vordhosbn");
932 assert_eq!(row.artist_name, "Aphex Twin");
933 assert_eq!(row.album_title, "Drukqs");
934 assert_eq!(row.remote_id.as_deref(), Some("remote-001"));
935 assert_eq!(row.source, "remote");
936 }
937
938 struct StubServer {
943 addr: std::net::SocketAddr,
944 shutdown: Arc<std::sync::atomic::AtomicBool>,
945 }
946
947 #[derive(Default)]
948 struct StubState {
949 albums: Mutex<Vec<(String, String, String)>>,
951 failing: Mutex<HashSet<String>>,
953 unlisted: Mutex<HashSet<String>>,
956 insert_after_first_page: Mutex<Option<(String, String, String)>>,
959 list_pages_served: AtomicUsize,
960 list_types: Mutex<Vec<String>>,
961 album_calls: Mutex<Vec<String>>,
962 lists_songs: bool,
966 song_pages: Mutex<Vec<usize>>,
968 failing_song_pages: Mutex<HashSet<usize>>,
970 library_modified: Mutex<Option<i64>>,
972 }
973
974 impl StubServer {
975 fn start(state: Arc<StubState>) -> Self {
976 let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
977 listener.set_nonblocking(true).unwrap();
978 let addr = listener.local_addr().unwrap();
979 let shutdown = Arc::new(std::sync::atomic::AtomicBool::new(false));
980
981 let stop = shutdown.clone();
982 std::thread::spawn(move || {
983 while !stop.load(Ordering::Relaxed) {
984 match listener.accept() {
985 Ok((stream, _)) => {
986 let _ = stream.set_nonblocking(false);
988 let state = state.clone();
989 std::thread::spawn(move || handle(stream, state));
990 }
991 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
992 std::thread::sleep(std::time::Duration::from_millis(2));
993 }
994 Err(_) => break,
995 }
996 }
997 });
998
999 Self { addr, shutdown }
1000 }
1001
1002 fn url(&self) -> String {
1003 format!("http://{}", self.addr)
1004 }
1005 }
1006
1007 impl Drop for StubServer {
1008 fn drop(&mut self) {
1009 self.shutdown.store(true, Ordering::Relaxed);
1010 }
1011 }
1012
1013 fn handle(mut stream: std::net::TcpStream, state: Arc<StubState>) {
1017 use std::io::{BufRead, Write};
1018
1019 let Ok(peek) = stream.try_clone() else { return };
1020 let mut reader = std::io::BufReader::new(peek);
1021
1022 loop {
1023 let mut request_line = String::new();
1024 if reader.read_line(&mut request_line).unwrap_or(0) == 0 {
1025 return;
1026 }
1027 let mut line = String::new();
1028 while reader.read_line(&mut line).unwrap_or(0) > 0 {
1029 if line == "\r\n" || line == "\n" {
1030 break;
1031 }
1032 line.clear();
1033 }
1034
1035 let target = request_line.split_whitespace().nth(1).unwrap_or("/");
1036 let (path, query) = target.split_once('?').unwrap_or((target, ""));
1037 let params: std::collections::HashMap<&str, &str> = query
1038 .split('&')
1039 .filter_map(|kv| kv.split_once('='))
1040 .collect();
1041
1042 let (status, body) = respond(&state, path, ¶ms);
1043
1044 let write = write!(
1045 stream,
1046 "HTTP/1.1 {} OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n",
1047 status,
1048 body.len()
1049 )
1050 .and_then(|()| stream.write_all(body.as_bytes()))
1051 .and_then(|()| stream.flush());
1052 if write.is_err() {
1053 return;
1054 }
1055 }
1056 }
1057
1058 fn respond(
1059 state: &Arc<StubState>,
1060 path: &str,
1061 params: &std::collections::HashMap<&str, &str>,
1062 ) -> (u16, String) {
1063 match path.rsplit('/').next().unwrap_or("") {
1064 "getArtists" => (
1065 200,
1066 r#"{"subsonic-response":{"status":"ok","artists":{"index":[{"artist":[{"id":"ar1","name":"Stub Artist"}]}]}}}"#.to_string(),
1067 ),
1068 "getAlbumList2" => {
1069 state
1070 .list_types
1071 .lock()
1072 .unwrap()
1073 .push(params.get("type").copied().unwrap_or("").to_string());
1074 let offset: usize = params.get("offset").and_then(|o| o.parse().ok()).unwrap_or(0);
1075 let size: usize = params.get("size").and_then(|s| s.parse().ok()).unwrap_or(500);
1076
1077 let albums = state.albums.lock().unwrap();
1078 let unlisted = state.unlisted.lock().unwrap();
1079 let slice: Vec<String> = albums
1080 .iter()
1081 .filter(|(id, _, _)| !unlisted.contains(id))
1082 .skip(offset)
1083 .take(size)
1084 .map(|(id, name, created)| {
1085 format!(
1086 r#"{{"id":"{}","name":"{}","artist":"Stub Artist","created":"{}","songCount":1}}"#,
1087 id, name, created
1088 )
1089 })
1090 .collect();
1091 drop(albums);
1092 drop(unlisted);
1093
1094 if state.list_pages_served.fetch_add(1, Ordering::SeqCst) == 0
1095 && let Some(new_album) = state.insert_after_first_page.lock().unwrap().take()
1096 {
1097 state.albums.lock().unwrap().insert(0, new_album);
1098 }
1099
1100 (
1101 200,
1102 format!(
1103 r#"{{"subsonic-response":{{"status":"ok","albumList2":{{"album":[{}]}}}}}}"#,
1104 slice.join(",")
1105 ),
1106 )
1107 }
1108 "getIndexes" => match *state.library_modified.lock().unwrap() {
1109 Some(at) => (
1110 200,
1111 format!(
1112 r#"{{"subsonic-response":{{"status":"ok","indexes":{{"lastModified":{at}}}}}}}"#
1113 ),
1114 ),
1115 None => (200, r#"{"subsonic-response":{"status":"ok"}}"#.to_string()),
1116 },
1117 "search3" => {
1118 let offset: usize = params.get("songOffset").and_then(|o| o.parse().ok()).unwrap_or(0);
1119 let size: usize = params.get("songCount").and_then(|s| s.parse().ok()).unwrap_or(20);
1120 state.song_pages.lock().unwrap().push(offset);
1121 if state.failing_song_pages.lock().unwrap().contains(&offset) {
1122 return (500, r#"{"error":"boom"}"#.to_string());
1123 }
1124 let songs: Vec<String> = if state.lists_songs {
1125 state
1126 .albums
1127 .lock()
1128 .unwrap()
1129 .iter()
1130 .skip(offset)
1131 .take(size)
1132 .map(|(id, name, _)| {
1133 format!(
1134 r#"{{"id":"s{id}","title":"Song {id}","albumId":"{id}","album":"{name}","track":1,"suffix":"flac"}}"#
1135 )
1136 })
1137 .collect()
1138 } else {
1139 Vec::new()
1140 };
1141 (
1142 200,
1143 format!(
1144 r#"{{"subsonic-response":{{"status":"ok","searchResult3":{{"song":[{}]}}}}}}"#,
1145 songs.join(",")
1146 ),
1147 )
1148 }
1149 "getAlbum" => {
1150 let id = params.get("id").copied().unwrap_or("");
1151 state.album_calls.lock().unwrap().push(id.to_string());
1152 if state.failing.lock().unwrap().contains(id) {
1153 (500, r#"{"error":"boom"}"#.to_string())
1154 } else if !state.albums.lock().unwrap().iter().any(|(a, _, _)| a == id) {
1155 (
1156 200,
1157 r#"{"subsonic-response":{"status":"failed","error":{"code":70,"message":"not found"}}}"#
1158 .to_string(),
1159 )
1160 } else {
1161 (
1162 200,
1163 format!(
1164 r#"{{"subsonic-response":{{"status":"ok","album":{{"id":"{id}","name":"{name}","artist":"Stub Artist","song":[{{"id":"s{id}","title":"Song {id}","track":1,"suffix":"flac"}}]}}}}}}"#,
1165 name = state
1166 .albums
1167 .lock()
1168 .unwrap()
1169 .iter()
1170 .find(|(a, _, _)| a == id)
1171 .map_or_else(|| format!("Album {id}"), |(_, n, _)| n.clone())
1172 ),
1173 )
1174 }
1175 }
1176 _ => (404, "{}".to_string()),
1177 }
1178 }
1179
1180 fn stub_albums(n: usize) -> Vec<(String, String, String)> {
1181 (0..n)
1182 .map(|i| {
1183 (
1184 format!("a{:04}", i),
1185 format!("Album {:04}", i),
1186 "2024-01-15T10:30:00Z".to_string(),
1187 )
1188 })
1189 .collect()
1190 }
1191
1192 #[test]
1193 fn failed_album_fetch_does_not_advance_last_sync_and_next_sync_retries() {
1194 let (db, _dir) = test_db();
1195 let state = Arc::new(StubState {
1196 albums: Mutex::new(stub_albums(4)),
1197 failing: Mutex::new(["a0002".to_string()].into_iter().collect()),
1198 ..Default::default()
1199 });
1200 let server = StubServer::start(state.clone());
1201 let client = SubsonicClient::new(&server.url(), "u", "p");
1202
1203 let first = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1204 assert_eq!(first.albums_failed, 1, "the failing album must be counted");
1205 assert_eq!(first.albums_synced, 3);
1206 assert!(!first.is_complete());
1207 assert_eq!(
1208 get_last_sync(&db, &server.url()).unwrap(),
1209 None,
1210 "an incomplete sync must not advance last_sync"
1211 );
1212
1213 state.failing.lock().unwrap().clear();
1216 state.album_calls.lock().unwrap().clear();
1217 let second = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1218
1219 assert!(
1220 state
1221 .album_calls
1222 .lock()
1223 .unwrap()
1224 .contains(&"a0002".to_string()),
1225 "the previously failed album must be retried"
1226 );
1227 assert_eq!(second.albums_failed, 0);
1228 assert!(second.is_complete());
1229 assert!(
1230 get_last_sync(&db, &server.url()).unwrap().is_some(),
1231 "a clean sync advances last_sync"
1232 );
1233 }
1234
1235 #[test]
1236 fn album_inserted_mid_pagination_is_not_fetched_twice_or_skipped() {
1237 let (db, _dir) = test_db();
1241 let state = Arc::new(StubState {
1242 albums: Mutex::new(stub_albums(600)),
1243 insert_after_first_page: Mutex::new(Some((
1244 "aNEW".to_string(),
1245 "AAA Brand New".to_string(),
1246 "2024-06-01T00:00:00Z".to_string(),
1247 ))),
1248 ..Default::default()
1249 });
1250 let server = StubServer::start(state.clone());
1251 let client = SubsonicClient::new(&server.url(), "u", "p");
1252
1253 let result = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1254 assert_eq!(result.albums_failed, 0);
1255
1256 let calls = state.album_calls.lock().unwrap().clone();
1257 let unique: HashSet<&String> = calls.iter().collect();
1258 assert_eq!(
1259 calls.len(),
1260 unique.len(),
1261 "no album may be fetched twice after the window shifts"
1262 );
1263
1264 for i in 0..600 {
1266 let id = format!("a{:04}", i);
1267 assert!(unique.contains(&id), "album {} was skipped", id);
1268 }
1269
1270 let types = state.list_types.lock().unwrap().clone();
1271 assert!(
1272 types.iter().all(|t| t == "alphabeticalByName"),
1273 "the paginated walk must use a stable ordering, got {:?}",
1274 types
1275 );
1276 }
1277
1278 #[test]
1279 fn a_full_sync_relinks_files_whose_server_id_changed() {
1280 let (db, _dir) = test_db();
1281 let file = TrackMeta {
1282 date: None,
1283 disc: None,
1284 path: Some("/music/song.flac".into()),
1285 source: "local".into(),
1286 album_remote_id: None,
1287 artist_remote_id: None,
1288 mbid: None,
1289 ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1290 };
1291 queries::upsert_track(&db.conn, &file).unwrap();
1292
1293 let state = Arc::new(StubState {
1294 albums: Mutex::new(stub_albums(1)),
1295 ..Default::default()
1296 });
1297 let server = StubServer::start(state);
1298 let client = SubsonicClient::new(&server.url(), "u", "p");
1299 sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1300
1301 let rows: Vec<(Option<String>, Option<String>)> = db
1302 .conn
1303 .prepare("SELECT path, remote_id FROM tracks WHERE title = 'Song a0000'")
1304 .unwrap()
1305 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1306 .unwrap()
1307 .collect::<Result<_, _>>()
1308 .unwrap();
1309 assert_eq!(
1310 rows,
1311 vec![(Some("/music/song.flac".into()), Some("sa0000".into()))],
1312 "the file takes the server's current id, and the recording is listed once"
1313 );
1314 }
1315
1316 #[test]
1317 fn a_full_sync_folds_server_rows_whose_id_vanished() {
1318 let (db, _dir) = test_db();
1319 let ghost = TrackMeta {
1320 date: None,
1321 disc: None,
1322 album_remote_id: None,
1323 artist_remote_id: None,
1324 mbid: None,
1325 ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1326 };
1327 let ghost_id = queries::upsert_track(&db.conn, &ghost).unwrap();
1328 queries::add_favourite(&db.conn, queries::LOCAL_USER, ghost_id).unwrap();
1329
1330 let state = Arc::new(StubState {
1331 albums: Mutex::new(stub_albums(1)),
1332 ..Default::default()
1333 });
1334 let server = StubServer::start(state);
1335 let client = SubsonicClient::new(&server.url(), "u", "p");
1336 sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1337
1338 let rows: Vec<(i64, String)> = db
1339 .conn
1340 .prepare("SELECT id, remote_id FROM tracks WHERE title = 'Song a0000'")
1341 .unwrap()
1342 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1343 .unwrap()
1344 .collect::<Result<_, _>>()
1345 .unwrap();
1346 assert_eq!(
1347 rows,
1348 vec![(ghost_id, "sa0000".to_string())],
1349 "the dead row takes the live id"
1350 );
1351 assert_eq!(
1352 queries::load_favourites(&db.conn, queries::LOCAL_USER).unwrap(),
1353 HashSet::from([ghost_id]),
1354 "the favourite follows"
1355 );
1356 }
1357
1358 #[test]
1361 fn a_library_that_has_not_moved_is_not_walked_again() {
1362 use crate::helpers::{Walk, sync_remote};
1363 let (db, _dir) = test_db();
1364 let state = Arc::new(StubState {
1365 albums: Mutex::new(stub_albums(3)),
1366 library_modified: Mutex::new(Some(1_000)),
1367 ..Default::default()
1368 });
1369 let server = StubServer::start(state.clone());
1370 let client = SubsonicClient::new(&server.url(), "u", "p");
1371 let walks = || state.list_types.lock().unwrap().len();
1372 let sync = |walk| sync_remote(&db, &client, walk, &server.url(), "u", &|_| {}).unwrap();
1373
1374 sync(Walk::IfChanged);
1375 let first = walks();
1376 assert!(first > 0, "nothing recorded, so walked");
1377
1378 sync(Walk::IfChanged);
1379 assert_eq!(walks(), first, "unchanged, so not walked");
1380
1381 sync(Walk::Always);
1382 let asked = walks();
1383 assert!(asked > first, "asked for, so walked");
1384
1385 *state.library_modified.lock().unwrap() = Some(2_000);
1386 sync(Walk::IfChanged);
1387 assert!(walks() > asked, "moved, so walked");
1388 }
1389
1390 #[test]
1392 fn a_server_without_a_library_version_is_always_walked() {
1393 use crate::helpers::{Walk, sync_remote};
1394 let (db, _dir) = test_db();
1395 let state = Arc::new(StubState {
1396 albums: Mutex::new(stub_albums(3)),
1397 ..Default::default()
1398 });
1399 let server = StubServer::start(state.clone());
1400 let client = SubsonicClient::new(&server.url(), "u", "p");
1401 let walks = || state.list_types.lock().unwrap().len();
1402
1403 sync_remote(&db, &client, Walk::IfChanged, &server.url(), "u", &|_| {}).unwrap();
1404 let first = walks();
1405 sync_remote(&db, &client, Walk::IfChanged, &server.url(), "u", &|_| {}).unwrap();
1406 assert!(walks() > first);
1407 }
1408
1409 #[test]
1412 fn an_album_left_out_of_the_listing_is_removed_only_when_gone() {
1413 let (db, _dir) = test_db();
1414 let state = Arc::new(StubState {
1415 albums: Mutex::new(stub_albums(3)),
1416 ..Default::default()
1417 });
1418 let server = StubServer::start(state.clone());
1419 let client = SubsonicClient::new(&server.url(), "u", "p");
1420 sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1421
1422 state.unlisted.lock().unwrap().insert("a0001".into());
1423 let shifted = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1424 assert_eq!(shifted.tracks_removed, 0, "still there, only unlisted");
1425
1426 state
1427 .albums
1428 .lock()
1429 .unwrap()
1430 .retain(|(id, _, _)| id != "a0001");
1431 let deleted = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1432 assert_eq!(deleted.tracks_removed, 1);
1433 }
1434
1435 #[test]
1438 fn a_full_sync_pages_songs_instead_of_fetching_each_album() {
1439 let (db, _dir) = test_db();
1440 let state = Arc::new(StubState {
1441 albums: Mutex::new(stub_albums(1_234)),
1442 lists_songs: true,
1443 ..Default::default()
1444 });
1445 let server = StubServer::start(state.clone());
1446 let client = SubsonicClient::new(&server.url(), "u", "p");
1447
1448 let seen = Mutex::new(Vec::new());
1449 let result = sync_library(&db, &client, &server.url(), "u", &|p| {
1450 seen.lock().unwrap().push(p)
1451 })
1452 .unwrap();
1453
1454 assert!(state.album_calls.lock().unwrap().is_empty(), "no getAlbum");
1455 let mut pages = state.song_pages.lock().unwrap().clone();
1456 pages.sort();
1457 pages.dedup();
1458 assert_eq!(pages[..3], [0, 500, 1000]);
1459 assert_eq!(result.tracks_synced, 1_234);
1460 assert_eq!(result.albums_synced, 1_234);
1461 assert!(result.is_complete());
1462 assert!(get_last_sync(&db, &server.url()).unwrap().is_some());
1463
1464 let album: (Option<String>, Option<i32>) = db
1465 .conn
1466 .query_row(
1467 "SELECT al.remote_id, al.total_tracks FROM albums al WHERE al.title = 'Album 0007'",
1468 [],
1469 |r| Ok((r.get(0)?, r.get(1)?)),
1470 )
1471 .unwrap();
1472 assert_eq!(
1473 album,
1474 (Some("a0007".into()), Some(1)),
1475 "album metadata comes from the list"
1476 );
1477
1478 let seen = seen.into_inner().unwrap();
1479 let tracks: Vec<&SyncProgress> = seen
1480 .iter()
1481 .filter(|p| p.phase == SyncPhase::Tracks)
1482 .collect();
1483 assert!(tracks.len() >= 3, "progress at least per page");
1484 assert!(tracks.iter().all(|p| p.total == Some(1_234)));
1485 assert_eq!(tracks.last().unwrap().done, 1_234);
1486 assert_eq!(seen.last().unwrap().phase, SyncPhase::Finishing);
1487 }
1488
1489 #[test]
1492 fn a_server_that_lists_no_songs_is_synced_album_by_album() {
1493 let (db, _dir) = test_db();
1494 let state = Arc::new(StubState {
1495 albums: Mutex::new(stub_albums(3)),
1496 ..Default::default()
1497 });
1498 let server = StubServer::start(state.clone());
1499 let client = SubsonicClient::new(&server.url(), "u", "p");
1500
1501 let result = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1502
1503 assert_eq!(state.song_pages.lock().unwrap().clone(), vec![0]);
1504 assert_eq!(state.album_calls.lock().unwrap().len(), 3);
1505 assert_eq!(result.tracks_synced, 3);
1506 assert!(result.is_complete());
1507 }
1508
1509 #[test]
1512 fn a_failed_song_page_does_not_advance_last_sync() {
1513 let (db, _dir) = test_db();
1514 let state = Arc::new(StubState {
1515 albums: Mutex::new(stub_albums(1_200)),
1516 lists_songs: true,
1517 failing_song_pages: Mutex::new([500].into_iter().collect()),
1518 ..Default::default()
1519 });
1520 let server = StubServer::start(state.clone());
1521 let client = SubsonicClient::new(&server.url(), "u", "p");
1522
1523 let result = sync_library(&db, &client, &server.url(), "u", &|_| {}).unwrap();
1524
1525 assert_eq!(result.pages_failed, 1);
1526 assert!(!result.is_complete());
1527 assert_eq!(get_last_sync(&db, &server.url()).unwrap(), None);
1528 let retries = state
1529 .song_pages
1530 .lock()
1531 .unwrap()
1532 .iter()
1533 .filter(|&&o| o == 500)
1534 .count();
1535 assert_eq!(retries as u32, PAGE_ATTEMPTS);
1536 }
1537
1538 #[test]
1539 fn sync_deduplicates_by_remote_id() {
1540 let (db, _dir) = test_db();
1541
1542 let meta1 = remote_track_meta("remote-dup", "Original Title", "Artist A", "Album X");
1544 let id1 = queries::upsert_track(&db.conn, &meta1).unwrap();
1545
1546 let meta2 = remote_track_meta("remote-dup", "Updated Title", "Artist A", "Album X");
1548 let id2 = queries::upsert_track(&db.conn, &meta2).unwrap();
1549
1550 assert_eq!(id1, id2, "same remote_id should resolve to same track row");
1552
1553 let row = queries::get_track_row(&db.conn, id2)
1555 .unwrap()
1556 .expect("track should exist");
1557 assert_eq!(row.title, "Updated Title");
1558 assert_eq!(row.remote_id.as_deref(), Some("remote-dup"));
1559
1560 let stats = queries::library_stats(&db.conn).unwrap();
1562 assert_eq!(
1563 stats.total_tracks, 1,
1564 "should have exactly 1 track after dedup"
1565 );
1566 }
1567
1568 #[test]
1572 fn resyncing_a_track_under_a_new_server_id_keeps_one_row() {
1573 let (db, _dir) = test_db();
1574 let before = remote_track_meta("46215", "Archangel", "Burial", "Untrue");
1575 let row = queries::upsert_synced_track(&db.conn, &before, &HashSet::from(["46215".into()]))
1576 .unwrap();
1577 queries::add_favourite(&db.conn, queries::LOCAL_USER, row).unwrap();
1578
1579 let uid = "0199a0b2-7c4e-7d3a-9f1b-2c3d4e5f6a7b";
1580 let after = remote_track_meta(uid, "Archangel", "Burial", "Untrue");
1581 let again =
1582 queries::upsert_synced_track(&db.conn, &after, &HashSet::from([uid.into()])).unwrap();
1583
1584 assert_eq!(again, row, "the same row, not a second copy");
1585 assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 1);
1586 let (remote_id, album, artist): (String, String, String) = db
1587 .conn
1588 .query_row(
1589 "SELECT t.remote_id, al.remote_id, ar.remote_id FROM tracks t
1590 JOIN albums al ON al.id = t.album_id JOIN artists ar ON ar.id = t.artist_id",
1591 [],
1592 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1593 )
1594 .unwrap();
1595 assert_eq!(remote_id, uid);
1596 assert_eq!(album, format!("album-of-{uid}"));
1597 assert_eq!(artist, format!("artist-of-{uid}"));
1598 let adopted: String = db
1599 .conn
1600 .query_row("SELECT uid FROM tracks WHERE id = ?1", [row], |r| r.get(0))
1601 .unwrap();
1602 assert_eq!(adopted, uid, "the row takes the server's uid");
1603 let favourites = queries::load_favourites(&db.conn, queries::LOCAL_USER).unwrap();
1604 assert_eq!(
1605 favourites,
1606 HashSet::from([row]),
1607 "the favourite stays with the row"
1608 );
1609 }
1610
1611 #[test]
1614 fn identical_entries_in_one_sync_stay_two_rows() {
1615 let (db, _dir) = test_db();
1616 let mut seen = HashSet::new();
1617 for id in ["dup-1", "dup-2"] {
1618 seen.insert(id.to_string());
1619 queries::upsert_synced_track(
1620 &db.conn,
1621 &remote_track_meta(id, "Archangel", "Burial", "Untrue"),
1622 &seen,
1623 )
1624 .unwrap();
1625 }
1626 assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 2);
1627 }
1628}