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, 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_removed: usize,
35}
36
37impl SyncResult {
38 pub fn is_complete(&self) -> bool {
40 self.albums_failed == 0 && self.pages_failed == 0
41 }
42}
43
44pub fn get_last_sync(
46 db: &Database,
47 url: &str,
48) -> Result<Option<i64>, crate::db::connection::DbError> {
49 let result = db.conn.query_row(
50 "SELECT last_sync FROM remote_servers WHERE url = ?1",
51 params![url],
52 |row| row.get::<_, Option<i64>>(0),
53 );
54 match result {
55 Ok(ts) => Ok(ts),
56 Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
57 Err(e) => Err(e.into()),
58 }
59}
60
61pub fn update_last_sync(
63 db: &Database,
64 url: &str,
65 username: &str,
66 timestamp: i64,
67) -> Result<(), crate::db::connection::DbError> {
68 db.conn.execute(
69 "INSERT INTO remote_servers (url, username, last_sync)
70 VALUES (?1, ?2, ?3)
71 ON CONFLICT(url) DO UPDATE SET last_sync = ?3",
72 params![url, username, timestamp],
73 )?;
74 Ok(())
75}
76
77fn parse_iso8601_to_unix(s: &str) -> Option<i64> {
85 use chrono::{DateTime, FixedOffset, NaiveDateTime};
86
87 if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
89 return Some(dt.timestamp());
90 }
91
92 if let Ok(naive) = NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S%.f") {
95 return Some(naive.and_utc().timestamp());
96 }
97 if let Ok(naive) = NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S") {
98 return Some(naive.and_utc().timestamp());
99 }
100
101 if let Ok(dt) = DateTime::<FixedOffset>::parse_from_str(s, "%Y-%m-%d %H:%M:%S%:z") {
103 return Some(dt.timestamp());
104 }
105 if let Ok(naive) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S") {
106 return Some(naive.and_utc().timestamp());
107 }
108
109 None
110}
111
112#[derive(Debug, Clone, Copy, PartialEq, Eq)]
114pub enum SyncPhase {
115 Albums,
117 Tracks,
119 Artists,
121 Finishing,
123}
124
125#[derive(Debug, Clone, Copy, PartialEq, Eq)]
129pub struct SyncProgress {
130 pub phase: SyncPhase,
131 pub done: u64,
132 pub total: Option<u64>,
133}
134
135const PAGE_SIZE: u32 = 500;
137
138const FETCH_LANES: usize = 4;
142
143const PAGE_ATTEMPTS: u32 = 3;
145
146const PAGED_ABOVE: usize = 200;
150
151pub fn sync_library(
172 db: &Database,
173 client: &SubsonicClient,
174 full: bool,
175 server_url: &str,
176 username: &str,
177 progress: &(dyn Fn(SyncProgress) + Sync),
178) -> Result<SyncResult, SyncError> {
179 let mut result = SyncResult::default();
180
181 let last_sync = if full {
182 None
183 } else {
184 get_last_sync(db, server_url)?
185 };
186
187 match last_sync {
188 Some(ts) => log::info!("incremental sync (albums created after {})", ts),
189 None => log::info!("full sync"),
190 }
191
192 let sync_start = std::time::SystemTime::now()
193 .duration_since(std::time::UNIX_EPOCH)
194 .unwrap_or_default()
195 .as_secs() as i64;
196
197 let albums = list_albums(client, progress)?;
198
199 let held = last_sync.and_then(|_| held_albums(db));
205 let wanted: Vec<&SubsonicAlbum> = albums
206 .iter()
207 .filter(|a| match last_sync {
208 None => true,
209 Some(ts) => {
210 a.created
211 .as_deref()
212 .and_then(parse_iso8601_to_unix)
213 .is_none_or(|created| created >= ts)
214 || held.as_ref().is_some_and(|h| differs(a, h.get(&a.id)))
215 }
216 })
217 .collect();
218 let expected: u64 = wanted
219 .iter()
220 .filter_map(|a| a.song_count)
221 .map(|n| n.max(0) as u64)
222 .sum();
223
224 let mut song_ids: HashSet<String> = HashSet::new();
225 let mut total = (expected > 0).then_some(expected);
226 let walked = if (last_sync.is_none() || wanted.len() > PAGED_ABOVE) && !wanted.is_empty() {
227 if total.is_none() {
228 total = client.song_count().ok().flatten();
229 }
230 sync_all_songs(
231 db,
232 client,
233 &albums,
234 total,
235 &mut result,
236 &mut song_ids,
237 progress,
238 )?
239 } else {
240 false
241 };
242 if !walked {
243 sync_by_album(
244 db,
245 client,
246 &wanted,
247 total,
248 &mut result,
249 &mut song_ids,
250 progress,
251 )?;
252 }
253
254 let artists = client.get_artists()?;
258 result.artists_synced = artists.len();
259 progress(SyncProgress {
260 phase: SyncPhase::Artists,
261 done: 0,
262 total: Some(artists.len() as u64),
263 });
264 write_artists(db, &artists, &mut result);
265
266 progress(SyncProgress {
267 phase: SyncPhase::Finishing,
268 done: 0,
269 total: None,
270 });
271
272 let listed_everything = total.is_none_or(|n| song_ids.len() as u64 >= n);
277 if full && result.is_complete() && !song_ids.is_empty() && listed_everything {
278 match queries::relink_vanished_remote_ids(&db.conn, &song_ids) {
279 Ok(0) => {}
280 Ok(n) => log::info!("{n} files had ids the server no longer knows; relinked"),
281 Err(e) => log::warn!("failed to relink tracks with vanished remote ids: {e}"),
282 }
283 }
284
285 let live_albums: std::collections::HashSet<String> =
291 albums.iter().map(|a| a.id.clone()).collect();
292 let live_tracks = (full && result.is_complete() && !song_ids.is_empty() && listed_everything)
293 .then_some(&song_ids);
294 if result.is_complete() && !live_albums.is_empty() {
295 match queries::remove_vanished_remote(&db.conn, live_tracks, Some(&live_albums)) {
296 Ok(0) => {}
297 Ok(n) => {
298 result.tracks_removed = n;
299 log::info!("{n} tracks the server no longer has; removed");
300 }
301 Err(e) => log::warn!("failed to remove tracks the server deleted: {e}"),
302 }
303 }
304
305 if result.is_complete() {
306 update_last_sync(db, server_url, username, sync_start)?;
307 } else {
308 log::warn!(
309 "{} album(s) and {} page(s) failed to fetch — leaving last_sync unchanged so the next sync retries them",
310 result.albums_failed,
311 result.pages_failed,
312 );
313 }
314
315 log::info!(
316 "sync complete: {} artists, {} albums, {} tracks, {} failed",
317 result.artists_synced,
318 result.albums_synced,
319 result.tracks_synced,
320 result.albums_failed,
321 );
322
323 db.optimize();
324
325 Ok(result)
326}
327
328struct HeldAlbum {
331 title: String,
332 artist: String,
333 tracks: i64,
334 seconds: i64,
335}
336
337fn held_albums(db: &Database) -> Option<HashMap<String, HeldAlbum>> {
341 let read = || -> rusqlite::Result<HashMap<String, HeldAlbum>> {
342 let mut stmt = db.conn.prepare(
343 "SELECT al.remote_id, al.title, COALESCE(ar.name, ''),
344 COUNT(t.id), COALESCE(SUM(t.duration_ms), 0) / 1000
345 FROM albums al
346 LEFT JOIN artists ar ON ar.id = al.artist_id
347 LEFT JOIN tracks t ON t.album_id = al.id AND t.remote_id IS NOT NULL
348 WHERE al.remote_id IS NOT NULL
349 GROUP BY al.id",
350 )?;
351 stmt.query_map([], |r| {
352 Ok((
353 r.get::<_, String>(0)?,
354 HeldAlbum {
355 title: r.get(1)?,
356 artist: r.get(2)?,
357 tracks: r.get(3)?,
358 seconds: r.get(4)?,
359 },
360 ))
361 })?
362 .collect()
363 };
364 read()
365 .inspect_err(|e| log::warn!("could not read held albums; syncing by date alone: {e}"))
366 .ok()
367}
368
369fn differs(listed: &SubsonicAlbum, held: Option<&HeldAlbum>) -> bool {
372 let Some(held) = held else { return true };
373 listed.name != held.title
374 || listed.artist.as_deref().is_some_and(|a| a != held.artist)
375 || listed
376 .song_count
377 .is_some_and(|n| i64::from(n) != held.tracks)
378 || listed
379 .duration
380 .is_some_and(|d| (d - held.seconds).abs() > 2)
381}
382
383fn list_albums(
385 client: &SubsonicClient,
386 progress: &(dyn Fn(SyncProgress) + Sync),
387) -> Result<Vec<SubsonicAlbum>, SyncError> {
388 let mut albums = Vec::new();
389 let mut seen: HashSet<String> = HashSet::new();
392 let mut offset = 0u32;
393 loop {
394 let page = client.get_album_list("alphabeticalByName", PAGE_SIZE, offset)?;
395 let count = page.len() as u32;
396 offset += count;
397 albums.extend(page.into_iter().filter(|a| seen.insert(a.id.clone())));
398 progress(SyncProgress {
399 phase: SyncPhase::Albums,
400 done: albums.len() as u64,
401 total: None,
402 });
403 if count < PAGE_SIZE {
404 return Ok(albums);
405 }
406 }
407}
408
409fn sync_all_songs(
416 db: &Database,
417 client: &SubsonicClient,
418 albums: &[SubsonicAlbum],
419 total: Option<u64>,
420 result: &mut SyncResult,
421 song_ids: &mut HashSet<String>,
422 progress: &(dyn Fn(SyncProgress) + Sync),
423) -> Result<bool, SyncError> {
424 let first = match client.all_songs_page(PAGE_SIZE, 0) {
425 Ok(songs) if !songs.is_empty() => songs,
426 Ok(_) => {
427 log::info!("server lists no songs for an empty search; syncing album by album");
428 return Ok(false);
429 }
430 Err(e) => {
431 log::info!("server refused an empty search ({e}); syncing album by album");
432 return Ok(false);
433 }
434 };
435
436 let by_id: HashMap<&str, &SubsonicAlbum> = albums.iter().map(|a| (a.id.as_str(), a)).collect();
437 let mut albums_seen: HashSet<String> = HashSet::new();
438 let mut write = |songs: Vec<SubsonicSong>,
439 result: &mut SyncResult,
440 song_ids: &mut HashSet<String>|
441 -> Result<(), SyncError> {
442 let batch = group_by_album(songs, &by_id);
443 albums_seen.extend(batch.iter().map(|a| a.id.clone()));
444 write_albums(db, client, &batch, result, song_ids)?;
445 progress(SyncProgress {
446 phase: SyncPhase::Tracks,
447 done: song_ids.len() as u64,
448 total,
449 });
450 Ok(())
451 };
452
453 let short = (first.len() as u32) < PAGE_SIZE;
454 write(first, result, song_ids)?;
455
456 if !short {
457 let next = AtomicU32::new(PAGE_SIZE);
463 let end = AtomicU32::new(u32::MAX);
466 let (tx, rx) = std::sync::mpsc::sync_channel(FETCH_LANES);
467 let failed = std::thread::scope(|scope| -> Result<usize, SyncError> {
468 for _ in 0..FETCH_LANES {
469 let tx = tx.clone();
470 let (next, end) = (&next, &end);
471 scope.spawn(move || {
472 loop {
473 let offset = next.fetch_add(PAGE_SIZE, Ordering::Relaxed);
474 if offset >= end.load(Ordering::Relaxed) {
475 return;
476 }
477 let page = fetch_page(client, offset);
478 if page
483 .as_ref()
484 .map_or(true, |songs| (songs.len() as u32) < PAGE_SIZE)
485 {
486 end.fetch_min(offset, Ordering::Relaxed);
487 }
488 if tx.send((offset, page)).is_err() {
489 return;
490 }
491 }
492 });
493 }
494 drop(tx);
495
496 let mut failed = 0;
497 for (offset, page) in rx {
498 match page {
499 Ok(songs) => write(songs, result, song_ids)?,
500 Err(e) => {
501 log::warn!("failed to fetch songs from offset {offset}: {e}");
502 failed += 1;
503 }
504 }
505 }
506 Ok(failed)
507 })?;
508 result.pages_failed += failed;
509 }
510
511 result.albums_synced += albums_seen.len();
512 Ok(true)
513}
514
515fn fetch_page(
518 client: &SubsonicClient,
519 offset: u32,
520) -> Result<Vec<SubsonicSong>, super::client::SubsonicError> {
521 let mut attempt = 1;
522 loop {
523 match client.all_songs_page(PAGE_SIZE, offset) {
524 Ok(songs) => return Ok(songs),
525 Err(e) if attempt >= PAGE_ATTEMPTS => return Err(e),
526 Err(e) => {
527 log::debug!("songs from offset {offset}, attempt {attempt}: {e}");
528 std::thread::sleep(std::time::Duration::from_millis(250 * attempt as u64));
529 attempt += 1;
530 }
531 }
532 }
533}
534
535fn group_by_album(
541 songs: Vec<SubsonicSong>,
542 albums: &HashMap<&str, &SubsonicAlbum>,
543) -> Vec<SubsonicAlbumFull> {
544 let mut grouped: Vec<SubsonicAlbumFull> = Vec::new();
545 let mut index: HashMap<String, usize> = HashMap::new();
546 for song in songs {
547 let Some(album_id) = song.album_id.clone() else {
548 log::warn!("song {} has no album id; skipped", song.id);
549 continue;
550 };
551 let i = *index.entry(album_id.clone()).or_insert_with(|| {
552 grouped.push(match albums.get(album_id.as_str()) {
553 Some(album) => SubsonicAlbumFull {
554 id: album.id.clone(),
555 name: album.name.clone(),
556 artist: album.artist.clone(),
557 artist_id: album.artist_id.clone(),
558 year: album.year,
559 genre: album.genre.clone(),
560 song_count: album.song_count,
561 created: album.created.clone(),
562 music_brainz_id: album.music_brainz_id.clone(),
563 sort_name: album.sort_name.clone(),
564 record_labels: album.record_labels.clone(),
565 song: Vec::new(),
566 },
567 None => SubsonicAlbumFull {
568 id: album_id,
569 name: song.album.clone().unwrap_or_default(),
570 artist: song.artist.clone(),
571 artist_id: song.artist_id.clone(),
572 year: song.year,
573 genre: song.genre.clone(),
574 song_count: None,
575 created: None,
576 music_brainz_id: None,
577 sort_name: None,
578 record_labels: Vec::new(),
579 song: Vec::new(),
580 },
581 });
582 grouped.len() - 1
583 });
584 grouped[i].song.push(song);
585 }
586 grouped
587}
588
589const ALBUM_BATCH: usize = 100;
592
593fn sync_by_album(
597 db: &Database,
598 client: &SubsonicClient,
599 albums: &[&SubsonicAlbum],
600 total: Option<u64>,
601 result: &mut SyncResult,
602 song_ids: &mut HashSet<String>,
603 progress: &(dyn Fn(SyncProgress) + Sync),
604) -> Result<(), SyncError> {
605 for batch in albums.chunks(ALBUM_BATCH) {
606 let failures = AtomicUsize::new(0);
607 let fetched: Vec<SubsonicAlbumFull> = batch
608 .par_iter()
609 .filter_map(|album| match client.get_album(&album.id) {
610 Ok(full) => Some(full),
611 Err(e) => {
612 log::warn!("failed to fetch album {}: {}", album.id, e);
613 failures.fetch_add(1, Ordering::Relaxed);
614 None
615 }
616 })
617 .collect();
618 result.albums_failed += failures.into_inner();
619 result.albums_synced += fetched.len();
620
621 write_albums(db, client, &fetched, result, song_ids)?;
622 progress(SyncProgress {
623 phase: SyncPhase::Tracks,
624 done: song_ids.len() as u64,
625 total,
626 });
627 log::info!(
628 "synced {} albums ({} tracks) so far...",
629 result.albums_synced,
630 result.tracks_synced
631 );
632 }
633 Ok(())
634}
635
636fn write_artists(db: &Database, artists: &[SubsonicArtist], result: &mut SyncResult) {
641 if db.conn.execute_batch("BEGIN").is_err() {
642 return;
643 }
644 let mut enriched = 0;
645 for artist in artists {
646 match queries::enrich_remote_artist(
647 &db.conn,
648 &artist.id,
649 artist.music_brainz_id.as_deref(),
650 artist.sort_name.as_deref(),
651 ) {
652 Ok(()) => enriched += 1,
653 Err(e) => log::warn!(
654 "failed to record artist metadata for {}: {}",
655 artist.name,
656 e
657 ),
658 }
659 }
660 if db.conn.execute_batch("COMMIT").is_err() {
661 let _ = db.conn.execute_batch("ROLLBACK");
662 return;
663 }
664 result.artists_synced = enriched;
665 log::info!("recorded metadata for {enriched} artists");
666}
667
668fn write_albums(
670 db: &Database,
671 client: &SubsonicClient,
672 albums: &[SubsonicAlbumFull],
673 result: &mut SyncResult,
674 song_ids: &mut HashSet<String>,
675) -> Result<(), SyncError> {
676 db.conn
677 .execute_batch("BEGIN")
678 .map_err(crate::db::connection::DbError::from)?;
679
680 for album in albums {
681 let artist_name = album.artist.as_deref().unwrap_or("Unknown Artist");
682
683 for song in &album.song {
684 song_ids.insert(song.id.clone());
685 let meta = TrackMeta {
686 title: song.title.clone(),
687 artist: song
688 .artist
689 .clone()
690 .unwrap_or_else(|| artist_name.to_string()),
691 album_artist: album.artist.clone(),
692 album: album.name.clone(),
693 date: album.year.map(|y| y.to_string()),
694 disc: song.disc_number,
695 track_number: song.track,
696 genre: song.genre.clone().or_else(|| album.genre.clone()),
697 label: None,
698 duration_ms: song.duration.map(|d| d * 1000),
699 codec: song.suffix.clone(),
700 sample_rate: positive(song.sampling_rate),
708 bit_depth: positive(song.bit_depth),
709 channels: positive(song.channel_count),
710 bitrate: song.bit_rate,
711 size_bytes: None,
712 mtime: None,
713 path: None,
714 source: "remote".to_string(),
715 remote_id: Some(song.id.clone()),
716 remote_url: Some(client.stream_url_template(&song.id)),
717 album_remote_id: Some(album.id.clone()),
718 artist_remote_id: album.artist_id.clone(),
719 mbid: song.music_brainz_id.clone(),
720 album_mbid: album.music_brainz_id.clone(),
721 album_added_at: album.created.clone(),
722 };
723
724 match queries::upsert_synced_track(&db.conn, &meta, song_ids) {
725 Ok(_) => result.tracks_synced += 1,
726 Err(e) => log::warn!("failed to insert remote track {}: {}", song.title, e),
727 }
728 }
729
730 if let Err(e) = queries::enrich_remote_album(
734 &db.conn,
735 &album.id,
736 album.music_brainz_id.as_deref(),
737 album.sort_name.as_deref(),
738 album.song_count,
739 album
740 .record_labels
741 .iter()
742 .map(|l| l.name.as_str())
743 .find(|l| !l.is_empty()),
744 ) {
745 log::warn!("failed to record album metadata for {}: {}", album.name, e);
746 }
747 }
748
749 db.conn
750 .execute_batch("COMMIT")
751 .map_err(crate::db::connection::DbError::from)?;
752
753 Ok(())
754}
755
756fn positive(value: Option<i32>) -> Option<i32> {
759 value.filter(|v| *v > 0)
760}
761
762#[cfg(test)]
763mod tests {
764 use super::*;
765
766 #[test]
767 fn parse_rfc3339_with_z() {
768 assert_eq!(
770 parse_iso8601_to_unix("2024-01-15T10:30:00Z"),
771 Some(1705314600)
772 );
773 }
774
775 #[test]
776 fn parse_rfc3339_with_offset() {
777 assert_eq!(
779 parse_iso8601_to_unix("2024-01-15T10:30:00+05:30"),
780 Some(1705294800)
781 );
782 }
783
784 #[test]
785 fn parse_rfc3339_negative_offset() {
786 assert_eq!(
788 parse_iso8601_to_unix("2024-01-15T10:30:00-05:00"),
789 Some(1705332600)
790 );
791 }
792
793 #[test]
794 fn parse_fractional_seconds_z() {
795 assert_eq!(
796 parse_iso8601_to_unix("2024-01-15T10:30:00.123Z"),
797 Some(1705314600)
798 );
799 }
800
801 #[test]
802 fn parse_fractional_seconds_offset() {
803 assert_eq!(
804 parse_iso8601_to_unix("2024-01-15T10:30:00.999+00:00"),
805 Some(1705314600)
806 );
807 }
808
809 #[test]
810 fn parse_no_timezone_assumes_utc() {
811 assert_eq!(
812 parse_iso8601_to_unix("2024-01-15T10:30:00"),
813 Some(1705314600)
814 );
815 }
816
817 #[test]
818 fn parse_no_timezone_fractional() {
819 assert_eq!(
820 parse_iso8601_to_unix("2024-01-15T10:30:00.500"),
821 Some(1705314600)
822 );
823 }
824
825 #[test]
826 fn parse_space_separator_with_tz() {
827 assert_eq!(
828 parse_iso8601_to_unix("2024-01-15 10:30:00+00:00"),
829 Some(1705314600)
830 );
831 }
832
833 #[test]
834 fn parse_space_separator_no_tz() {
835 assert_eq!(
836 parse_iso8601_to_unix("2024-01-15 10:30:00"),
837 Some(1705314600)
838 );
839 }
840
841 #[test]
842 fn parse_garbage_returns_none() {
843 assert_eq!(parse_iso8601_to_unix("not-a-date"), None);
844 assert_eq!(parse_iso8601_to_unix(""), None);
845 assert_eq!(parse_iso8601_to_unix("2024"), None);
846 }
847
848 #[test]
849 fn parse_epoch() {
850 assert_eq!(parse_iso8601_to_unix("1970-01-01T00:00:00Z"), Some(0));
851 }
852
853 use crate::db::connection::Database;
856 use crate::db::queries;
857 use std::sync::{Arc, Mutex};
858
859 fn test_db() -> (Database, tempfile::TempDir) {
860 let dir = tempfile::tempdir().unwrap();
861 let db = Database::open(&dir.path().join("sync_test.db")).unwrap();
862 (db, dir)
863 }
864
865 fn remote_track_meta(remote_id: &str, title: &str, artist: &str, album: &str) -> TrackMeta {
867 TrackMeta {
868 title: title.into(),
869 artist: artist.into(),
870 album_artist: Some(artist.into()),
871 album: album.into(),
872 date: Some("2024".into()),
873 disc: Some(1),
874 track_number: Some(1),
875 genre: Some("Electronic".into()),
876 label: None,
877 duration_ms: Some(240_000),
878 codec: Some("FLAC".into()),
879 sample_rate: Some(44100),
880 bit_depth: Some(16),
881 channels: Some(2),
882 bitrate: Some(1000),
883 size_bytes: None,
884 mtime: None,
885 path: None,
886 source: "remote".into(),
887 remote_id: Some(remote_id.into()),
888 remote_url: Some(format!("https://example.com/stream?id={}", remote_id)),
889 album_remote_id: Some(format!("album-of-{remote_id}")),
890 artist_remote_id: Some(format!("artist-of-{remote_id}")),
891 mbid: Some(format!("mbid-of-{remote_id}")),
892 album_mbid: None,
893 album_added_at: None,
894 }
895 }
896
897 #[test]
900 fn a_zero_quality_figure_is_treated_as_absent() {
901 assert_eq!(positive(Some(0)), None);
902 assert_eq!(positive(Some(16)), Some(16));
903 assert_eq!(positive(None), None);
904 }
905
906 #[test]
910 fn album_metadata_from_the_server_is_recorded() {
911 let (db, _dir) = test_db();
912
913 let meta = remote_track_meta("remote-300", "Anguish", "Sleep", "Volume One");
914 queries::upsert_track(&db.conn, &meta).unwrap();
915
916 queries::enrich_remote_album(
917 &db.conn,
918 "album-of-remote-300",
919 Some("mb-album-1"),
920 Some("volume one"),
921 Some(6),
922 Some("Off The Disk"),
923 )
924 .unwrap();
925
926 let row: (Option<String>, Option<String>, Option<i32>, Option<String>) = db
927 .conn
928 .query_row(
929 "SELECT mbid, sort_name, total_tracks, label FROM albums WHERE title = 'Volume One'",
930 [],
931 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
932 )
933 .unwrap();
934 assert_eq!(row.0.as_deref(), Some("mb-album-1"));
935 assert_eq!(row.1.as_deref(), Some("volume one"));
936 assert_eq!(row.2, Some(6));
937 assert_eq!(row.3.as_deref(), Some("Off The Disk"));
938 }
939
940 #[test]
943 fn server_metadata_does_not_overwrite_what_tags_said() {
944 let (db, _dir) = test_db();
945
946 let mut meta = remote_track_meta("remote-301", "Dopesmoker", "Sleep", "Dopesmoker");
947 meta.label = Some("From The Tags".into());
948 queries::upsert_track(&db.conn, &meta).unwrap();
949
950 queries::enrich_remote_album(
951 &db.conn,
952 "album-of-remote-301",
953 None,
954 None,
955 None,
956 Some("From The Server"),
957 )
958 .unwrap();
959
960 let label: Option<String> = db
961 .conn
962 .query_row(
963 "SELECT label FROM albums WHERE title = 'Dopesmoker'",
964 [],
965 |r| r.get(0),
966 )
967 .unwrap();
968 assert_eq!(label.as_deref(), Some("From The Tags"));
969 }
970
971 #[test]
974 fn artist_metadata_from_the_server_is_recorded() {
975 let (db, _dir) = test_db();
976
977 let meta = remote_track_meta("remote-302", "Holy Mountain", "Sleep", "Holy Mountain");
978 queries::upsert_track(&db.conn, &meta).unwrap();
979
980 queries::enrich_remote_artist(
981 &db.conn,
982 "artist-of-remote-302",
983 Some("mb-artist-1"),
984 Some("sleep"),
985 )
986 .unwrap();
987
988 let row: (Option<String>, Option<String>) = db
989 .conn
990 .query_row(
991 "SELECT mbid, sort_name FROM artists WHERE name = 'Sleep'",
992 [],
993 |r| Ok((r.get(0)?, r.get(1)?)),
994 )
995 .unwrap();
996 assert_eq!(row.0.as_deref(), Some("mb-artist-1"));
997 assert_eq!(row.1.as_deref(), Some("sleep"));
998 }
999
1000 #[test]
1002 fn a_synced_track_keeps_its_musicbrainz_id() {
1003 let (db, _dir) = test_db();
1004
1005 let meta = remote_track_meta("remote-303", "Aquarian", "Sleep", "Dopesmoker");
1006 let id = queries::upsert_track(&db.conn, &meta).unwrap();
1007
1008 let mbid: Option<String> = db
1009 .conn
1010 .query_row("SELECT mbid FROM tracks WHERE id = ?1", [id], |r| r.get(0))
1011 .unwrap();
1012 assert_eq!(mbid.as_deref(), Some("mbid-of-remote-303"));
1013 }
1014
1015 #[test]
1019 fn a_synced_track_keeps_its_quality_figures() {
1020 let (db, _dir) = test_db();
1021
1022 let meta = remote_track_meta("remote-200", "Anguish", "Sleep", "Volume One");
1023 let id = queries::upsert_track(&db.conn, &meta).unwrap();
1024
1025 let row = queries::get_track_row(&db.conn, id).unwrap().unwrap();
1026 assert_eq!(row.sample_rate, Some(44100));
1027 assert_eq!(row.bit_depth, Some(16));
1028 assert_eq!(row.channels, Some(2));
1029 }
1030
1031 #[test]
1035 fn a_sync_records_the_album_and_artist_ids_too() {
1036 let (db, _dir) = test_db();
1037
1038 let meta = remote_track_meta("remote-100", "Enter", "Russian Circles", "Enter");
1039 queries::upsert_track(&db.conn, &meta).unwrap();
1040
1041 let album: Option<String> = db
1042 .conn
1043 .query_row(
1044 "SELECT remote_id FROM albums WHERE title = 'Enter'",
1045 [],
1046 |r| r.get(0),
1047 )
1048 .unwrap();
1049 assert_eq!(album.as_deref(), Some("album-of-remote-100"));
1050
1051 let artist: Option<String> = db
1052 .conn
1053 .query_row(
1054 "SELECT remote_id FROM artists WHERE name = 'Russian Circles'",
1055 [],
1056 |r| r.get(0),
1057 )
1058 .unwrap();
1059 assert_eq!(artist.as_deref(), Some("artist-of-remote-100"));
1060 }
1061
1062 #[test]
1063 fn sync_upserts_tracks_to_database() {
1064 let (db, _dir) = test_db();
1065
1066 let meta = remote_track_meta("remote-001", "Vordhosbn", "Aphex Twin", "Drukqs");
1067 let track_id = queries::upsert_track(&db.conn, &meta).unwrap();
1068 assert!(track_id > 0, "upsert should return a valid track ID");
1069
1070 let row = queries::get_track_row(&db.conn, track_id)
1072 .unwrap()
1073 .expect("track should exist in DB");
1074 assert_eq!(row.title, "Vordhosbn");
1075 assert_eq!(row.artist_name, "Aphex Twin");
1076 assert_eq!(row.album_title, "Drukqs");
1077 assert_eq!(row.remote_id.as_deref(), Some("remote-001"));
1078 assert_eq!(row.source, "remote");
1079 }
1080
1081 struct StubServer {
1086 addr: std::net::SocketAddr,
1087 shutdown: Arc<std::sync::atomic::AtomicBool>,
1088 }
1089
1090 #[derive(Default)]
1091 struct StubState {
1092 albums: Mutex<Vec<(String, String, String)>>,
1094 failing: Mutex<HashSet<String>>,
1096 insert_after_first_page: Mutex<Option<(String, String, String)>>,
1099 list_pages_served: AtomicUsize,
1100 list_types: Mutex<Vec<String>>,
1101 album_calls: Mutex<Vec<String>>,
1102 lists_songs: bool,
1106 song_pages: Mutex<Vec<usize>>,
1108 failing_song_pages: Mutex<HashSet<usize>>,
1110 }
1111
1112 impl StubServer {
1113 fn start(state: Arc<StubState>) -> Self {
1114 let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
1115 listener.set_nonblocking(true).unwrap();
1116 let addr = listener.local_addr().unwrap();
1117 let shutdown = Arc::new(std::sync::atomic::AtomicBool::new(false));
1118
1119 let stop = shutdown.clone();
1120 std::thread::spawn(move || {
1121 while !stop.load(Ordering::Relaxed) {
1122 match listener.accept() {
1123 Ok((stream, _)) => {
1124 let _ = stream.set_nonblocking(false);
1126 let state = state.clone();
1127 std::thread::spawn(move || handle(stream, state));
1128 }
1129 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
1130 std::thread::sleep(std::time::Duration::from_millis(2));
1131 }
1132 Err(_) => break,
1133 }
1134 }
1135 });
1136
1137 Self { addr, shutdown }
1138 }
1139
1140 fn url(&self) -> String {
1141 format!("http://{}", self.addr)
1142 }
1143 }
1144
1145 impl Drop for StubServer {
1146 fn drop(&mut self) {
1147 self.shutdown.store(true, Ordering::Relaxed);
1148 }
1149 }
1150
1151 fn handle(mut stream: std::net::TcpStream, state: Arc<StubState>) {
1155 use std::io::{BufRead, Write};
1156
1157 let Ok(peek) = stream.try_clone() else { return };
1158 let mut reader = std::io::BufReader::new(peek);
1159
1160 loop {
1161 let mut request_line = String::new();
1162 if reader.read_line(&mut request_line).unwrap_or(0) == 0 {
1163 return;
1164 }
1165 let mut line = String::new();
1166 while reader.read_line(&mut line).unwrap_or(0) > 0 {
1167 if line == "\r\n" || line == "\n" {
1168 break;
1169 }
1170 line.clear();
1171 }
1172
1173 let target = request_line.split_whitespace().nth(1).unwrap_or("/");
1174 let (path, query) = target.split_once('?').unwrap_or((target, ""));
1175 let params: std::collections::HashMap<&str, &str> = query
1176 .split('&')
1177 .filter_map(|kv| kv.split_once('='))
1178 .collect();
1179
1180 let (status, body) = respond(&state, path, ¶ms);
1181
1182 let write = write!(
1183 stream,
1184 "HTTP/1.1 {} OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n",
1185 status,
1186 body.len()
1187 )
1188 .and_then(|()| stream.write_all(body.as_bytes()))
1189 .and_then(|()| stream.flush());
1190 if write.is_err() {
1191 return;
1192 }
1193 }
1194 }
1195
1196 fn respond(
1197 state: &Arc<StubState>,
1198 path: &str,
1199 params: &std::collections::HashMap<&str, &str>,
1200 ) -> (u16, String) {
1201 match path.rsplit('/').next().unwrap_or("") {
1202 "getArtists" => (
1203 200,
1204 r#"{"subsonic-response":{"status":"ok","artists":{"index":[{"artist":[{"id":"ar1","name":"Stub Artist"}]}]}}}"#.to_string(),
1205 ),
1206 "getAlbumList2" => {
1207 state
1208 .list_types
1209 .lock()
1210 .unwrap()
1211 .push(params.get("type").copied().unwrap_or("").to_string());
1212 let offset: usize = params.get("offset").and_then(|o| o.parse().ok()).unwrap_or(0);
1213 let size: usize = params.get("size").and_then(|s| s.parse().ok()).unwrap_or(500);
1214
1215 let albums = state.albums.lock().unwrap();
1216 let slice: Vec<String> = albums
1217 .iter()
1218 .skip(offset)
1219 .take(size)
1220 .map(|(id, name, created)| {
1221 format!(
1222 r#"{{"id":"{}","name":"{}","artist":"Stub Artist","created":"{}","songCount":1}}"#,
1223 id, name, created
1224 )
1225 })
1226 .collect();
1227 drop(albums);
1228
1229 if state.list_pages_served.fetch_add(1, Ordering::SeqCst) == 0
1230 && let Some(new_album) = state.insert_after_first_page.lock().unwrap().take()
1231 {
1232 state.albums.lock().unwrap().insert(0, new_album);
1233 }
1234
1235 (
1236 200,
1237 format!(
1238 r#"{{"subsonic-response":{{"status":"ok","albumList2":{{"album":[{}]}}}}}}"#,
1239 slice.join(",")
1240 ),
1241 )
1242 }
1243 "search3" => {
1244 let offset: usize = params.get("songOffset").and_then(|o| o.parse().ok()).unwrap_or(0);
1245 let size: usize = params.get("songCount").and_then(|s| s.parse().ok()).unwrap_or(20);
1246 state.song_pages.lock().unwrap().push(offset);
1247 if state.failing_song_pages.lock().unwrap().contains(&offset) {
1248 return (500, r#"{"error":"boom"}"#.to_string());
1249 }
1250 let songs: Vec<String> = if state.lists_songs {
1251 state
1252 .albums
1253 .lock()
1254 .unwrap()
1255 .iter()
1256 .skip(offset)
1257 .take(size)
1258 .map(|(id, name, _)| {
1259 format!(
1260 r#"{{"id":"s{id}","title":"Song {id}","albumId":"{id}","album":"{name}","track":1,"suffix":"flac"}}"#
1261 )
1262 })
1263 .collect()
1264 } else {
1265 Vec::new()
1266 };
1267 (
1268 200,
1269 format!(
1270 r#"{{"subsonic-response":{{"status":"ok","searchResult3":{{"song":[{}]}}}}}}"#,
1271 songs.join(",")
1272 ),
1273 )
1274 }
1275 "getAlbum" => {
1276 let id = params.get("id").copied().unwrap_or("");
1277 state.album_calls.lock().unwrap().push(id.to_string());
1278 if state.failing.lock().unwrap().contains(id) {
1279 (500, r#"{"error":"boom"}"#.to_string())
1280 } else {
1281 (
1282 200,
1283 format!(
1284 r#"{{"subsonic-response":{{"status":"ok","album":{{"id":"{id}","name":"{name}","artist":"Stub Artist","song":[{{"id":"s{id}","title":"Song {id}","track":1,"suffix":"flac"}}]}}}}}}"#,
1285 name = state
1286 .albums
1287 .lock()
1288 .unwrap()
1289 .iter()
1290 .find(|(a, _, _)| a == id)
1291 .map_or_else(|| format!("Album {id}"), |(_, n, _)| n.clone())
1292 ),
1293 )
1294 }
1295 }
1296 _ => (404, "{}".to_string()),
1297 }
1298 }
1299
1300 fn stub_albums(n: usize) -> Vec<(String, String, String)> {
1301 (0..n)
1302 .map(|i| {
1303 (
1304 format!("a{:04}", i),
1305 format!("Album {:04}", i),
1306 "2024-01-15T10:30:00Z".to_string(),
1307 )
1308 })
1309 .collect()
1310 }
1311
1312 #[test]
1313 fn failed_album_fetch_does_not_advance_last_sync_and_next_sync_retries() {
1314 let (db, _dir) = test_db();
1315 let state = Arc::new(StubState {
1316 albums: Mutex::new(stub_albums(4)),
1317 failing: Mutex::new(["a0002".to_string()].into_iter().collect()),
1318 ..Default::default()
1319 });
1320 let server = StubServer::start(state.clone());
1321 let client = SubsonicClient::new(&server.url(), "u", "p");
1322
1323 let first = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1324 assert_eq!(first.albums_failed, 1, "the failing album must be counted");
1325 assert_eq!(first.albums_synced, 3);
1326 assert!(!first.is_complete());
1327 assert_eq!(
1328 get_last_sync(&db, &server.url()).unwrap(),
1329 None,
1330 "an incomplete sync must not advance last_sync"
1331 );
1332
1333 state.failing.lock().unwrap().clear();
1336 state.album_calls.lock().unwrap().clear();
1337 let second = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1338
1339 assert!(
1340 state
1341 .album_calls
1342 .lock()
1343 .unwrap()
1344 .contains(&"a0002".to_string()),
1345 "the previously failed album must be retried"
1346 );
1347 assert_eq!(second.albums_failed, 0);
1348 assert!(second.is_complete());
1349 assert!(
1350 get_last_sync(&db, &server.url()).unwrap().is_some(),
1351 "a clean sync advances last_sync"
1352 );
1353 }
1354
1355 #[test]
1356 fn album_inserted_mid_pagination_is_not_fetched_twice_or_skipped() {
1357 let (db, _dir) = test_db();
1361 let state = Arc::new(StubState {
1362 albums: Mutex::new(stub_albums(600)),
1363 insert_after_first_page: Mutex::new(Some((
1364 "aNEW".to_string(),
1365 "AAA Brand New".to_string(),
1366 "2024-06-01T00:00:00Z".to_string(),
1367 ))),
1368 ..Default::default()
1369 });
1370 let server = StubServer::start(state.clone());
1371 let client = SubsonicClient::new(&server.url(), "u", "p");
1372
1373 let result = sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1374 assert_eq!(result.albums_failed, 0);
1375
1376 let calls = state.album_calls.lock().unwrap().clone();
1377 let unique: HashSet<&String> = calls.iter().collect();
1378 assert_eq!(
1379 calls.len(),
1380 unique.len(),
1381 "no album may be fetched twice after the window shifts"
1382 );
1383
1384 for i in 0..600 {
1386 let id = format!("a{:04}", i);
1387 assert!(unique.contains(&id), "album {} was skipped", id);
1388 }
1389
1390 let types = state.list_types.lock().unwrap().clone();
1391 assert!(
1392 types.iter().all(|t| t == "alphabeticalByName"),
1393 "the paginated walk must use a stable ordering, got {:?}",
1394 types
1395 );
1396 }
1397
1398 #[test]
1399 fn a_full_sync_relinks_files_whose_server_id_changed() {
1400 let (db, _dir) = test_db();
1401 let file = TrackMeta {
1402 date: None,
1403 disc: None,
1404 path: Some("/music/song.flac".into()),
1405 source: "local".into(),
1406 album_remote_id: None,
1407 artist_remote_id: None,
1408 mbid: None,
1409 ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1410 };
1411 queries::upsert_track(&db.conn, &file).unwrap();
1412
1413 let state = Arc::new(StubState {
1414 albums: Mutex::new(stub_albums(1)),
1415 ..Default::default()
1416 });
1417 let server = StubServer::start(state);
1418 let client = SubsonicClient::new(&server.url(), "u", "p");
1419 sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1420
1421 let rows: Vec<(Option<String>, Option<String>)> = db
1422 .conn
1423 .prepare("SELECT path, remote_id FROM tracks WHERE title = 'Song a0000'")
1424 .unwrap()
1425 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1426 .unwrap()
1427 .collect::<Result<_, _>>()
1428 .unwrap();
1429 assert_eq!(
1430 rows,
1431 vec![(Some("/music/song.flac".into()), Some("sa0000".into()))],
1432 "the file takes the server's current id, and the recording is listed once"
1433 );
1434 }
1435
1436 #[test]
1437 fn a_full_sync_folds_server_rows_whose_id_vanished() {
1438 let (db, _dir) = test_db();
1439 let ghost = TrackMeta {
1440 date: None,
1441 disc: None,
1442 album_remote_id: None,
1443 artist_remote_id: None,
1444 mbid: None,
1445 ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1446 };
1447 let ghost_id = queries::upsert_track(&db.conn, &ghost).unwrap();
1448 let ghost_url = ghost.remote_url.clone().unwrap();
1449 db.conn
1450 .execute(
1451 "INSERT INTO favourites (track_path) VALUES (?1)",
1452 params![ghost_url],
1453 )
1454 .unwrap();
1455
1456 let state = Arc::new(StubState {
1457 albums: Mutex::new(stub_albums(1)),
1458 ..Default::default()
1459 });
1460 let server = StubServer::start(state);
1461 let client = SubsonicClient::new(&server.url(), "u", "p");
1462 sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1463
1464 let rows: Vec<(i64, String, String)> = db
1465 .conn
1466 .prepare("SELECT id, remote_id, remote_url FROM tracks WHERE title = 'Song a0000'")
1467 .unwrap()
1468 .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))
1469 .unwrap()
1470 .collect::<Result<_, _>>()
1471 .unwrap();
1472 assert_eq!(rows.len(), 1, "the dead row takes the live id");
1473 let (id, remote_id, remote_url) = &rows[0];
1474 assert_eq!(*id, ghost_id);
1475 assert_eq!(remote_id, "sa0000");
1476 let favourites: Vec<String> = db
1477 .conn
1478 .prepare("SELECT track_path FROM favourites")
1479 .unwrap()
1480 .query_map([], |row| row.get(0))
1481 .unwrap()
1482 .collect::<Result<_, _>>()
1483 .unwrap();
1484 assert_eq!(
1485 favourites,
1486 vec![remote_url.clone()],
1487 "the favourite follows"
1488 );
1489 }
1490
1491 #[test]
1492 fn incremental_sync_only_fetches_albums_created_after_last_sync() {
1493 let (db, _dir) = test_db();
1494 let mut albums = stub_albums(3);
1495 albums[0].2 = "2020-01-01T00:00:00Z".into();
1496 albums[1].2 = "2020-01-01T00:00:00Z".into();
1497 albums[2].2 = "2030-01-01T00:00:00Z".into();
1498
1499 let state = Arc::new(StubState {
1500 albums: Mutex::new(albums),
1501 ..Default::default()
1502 });
1503 let server = StubServer::start(state.clone());
1504 let client = SubsonicClient::new(&server.url(), "u", "p");
1505
1506 sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1509 state.album_calls.lock().unwrap().clear();
1510 let watermark = parse_iso8601_to_unix("2025-01-01T00:00:00Z").unwrap();
1511 update_last_sync(&db, &server.url(), "u", watermark).unwrap();
1512
1513 let result = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1514
1515 assert_eq!(result.albums_synced, 1, "only the new album needs fetching");
1516 assert_eq!(
1517 *state.album_calls.lock().unwrap(),
1518 vec!["a0002".to_string()]
1519 );
1520 }
1521
1522 #[test]
1525 fn incremental_sync_fetches_an_album_the_server_now_lists_differently() {
1526 let (db, _dir) = test_db();
1527 let mut albums = stub_albums(3);
1528 for a in &mut albums {
1529 a.2 = "2020-01-01T00:00:00Z".into();
1530 }
1531 let state = Arc::new(StubState {
1532 albums: Mutex::new(albums),
1533 ..Default::default()
1534 });
1535 let server = StubServer::start(state.clone());
1536 let client = SubsonicClient::new(&server.url(), "u", "p");
1537 sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1538 state.album_calls.lock().unwrap().clear();
1539
1540 state.albums.lock().unwrap()[1].1 = "Album 0001 (Retitled)".into();
1541 let result = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1542
1543 assert_eq!(
1544 *state.album_calls.lock().unwrap(),
1545 vec!["a0001".to_string()]
1546 );
1547 assert_eq!(result.albums_synced, 1);
1548 let title: String = db
1549 .conn
1550 .query_row(
1551 "SELECT title FROM albums WHERE remote_id = 'a0001'",
1552 [],
1553 |r| r.get(0),
1554 )
1555 .unwrap();
1556 assert_eq!(title, "Album 0001 (Retitled)");
1557 }
1558
1559 #[test]
1562 fn incremental_sync_fetches_an_album_it_does_not_hold() {
1563 let (db, _dir) = test_db();
1564 let mut albums = stub_albums(2);
1565 for a in &mut albums {
1566 a.2 = "2020-01-01T00:00:00Z".into();
1567 }
1568 let state = Arc::new(StubState {
1569 albums: Mutex::new(albums),
1570 ..Default::default()
1571 });
1572 let server = StubServer::start(state.clone());
1573 let client = SubsonicClient::new(&server.url(), "u", "p");
1574 sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1575 state.album_calls.lock().unwrap().clear();
1576
1577 state.albums.lock().unwrap().push((
1578 "a0099".into(),
1579 "Album 0099".into(),
1580 "2020-01-01T00:00:00Z".into(),
1581 ));
1582 sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1583
1584 assert_eq!(
1585 *state.album_calls.lock().unwrap(),
1586 vec!["a0099".to_string()]
1587 );
1588 }
1589
1590 #[test]
1591 fn album_with_unparseable_created_is_always_fetched() {
1592 let (db, _dir) = test_db();
1593 let state = Arc::new(StubState {
1594 albums: Mutex::new(vec![("a0000".into(), "Album".into(), "who knows".into())]),
1595 ..Default::default()
1596 });
1597 let server = StubServer::start(state.clone());
1598 let client = SubsonicClient::new(&server.url(), "u", "p");
1599
1600 update_last_sync(&db, &server.url(), "u", 4_000_000_000).unwrap();
1601 let result = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1602
1603 assert_eq!(
1604 result.albums_synced, 1,
1605 "an album with no usable timestamp must not be assumed old"
1606 );
1607 }
1608
1609 #[test]
1612 fn a_full_sync_pages_songs_instead_of_fetching_each_album() {
1613 let (db, _dir) = test_db();
1614 let state = Arc::new(StubState {
1615 albums: Mutex::new(stub_albums(1_234)),
1616 lists_songs: true,
1617 ..Default::default()
1618 });
1619 let server = StubServer::start(state.clone());
1620 let client = SubsonicClient::new(&server.url(), "u", "p");
1621
1622 let seen = Mutex::new(Vec::new());
1623 let result = sync_library(&db, &client, true, &server.url(), "u", &|p| {
1624 seen.lock().unwrap().push(p)
1625 })
1626 .unwrap();
1627
1628 assert!(state.album_calls.lock().unwrap().is_empty(), "no getAlbum");
1629 let mut pages = state.song_pages.lock().unwrap().clone();
1630 pages.sort();
1631 pages.dedup();
1632 assert_eq!(pages[..3], [0, 500, 1000]);
1633 assert_eq!(result.tracks_synced, 1_234);
1634 assert_eq!(result.albums_synced, 1_234);
1635 assert!(result.is_complete());
1636 assert!(get_last_sync(&db, &server.url()).unwrap().is_some());
1637
1638 let album: (Option<String>, Option<i32>) = db
1639 .conn
1640 .query_row(
1641 "SELECT al.remote_id, al.total_tracks FROM albums al WHERE al.title = 'Album 0007'",
1642 [],
1643 |r| Ok((r.get(0)?, r.get(1)?)),
1644 )
1645 .unwrap();
1646 assert_eq!(
1647 album,
1648 (Some("a0007".into()), Some(1)),
1649 "album metadata comes from the list"
1650 );
1651
1652 let seen = seen.into_inner().unwrap();
1653 let tracks: Vec<&SyncProgress> = seen
1654 .iter()
1655 .filter(|p| p.phase == SyncPhase::Tracks)
1656 .collect();
1657 assert!(tracks.len() >= 3, "progress at least per page");
1658 assert!(tracks.iter().all(|p| p.total == Some(1_234)));
1659 assert_eq!(tracks.last().unwrap().done, 1_234);
1660 assert_eq!(seen.last().unwrap().phase, SyncPhase::Finishing);
1661 }
1662
1663 #[test]
1666 fn a_server_that_lists_no_songs_is_synced_album_by_album() {
1667 let (db, _dir) = test_db();
1668 let state = Arc::new(StubState {
1669 albums: Mutex::new(stub_albums(3)),
1670 ..Default::default()
1671 });
1672 let server = StubServer::start(state.clone());
1673 let client = SubsonicClient::new(&server.url(), "u", "p");
1674
1675 let result = sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1676
1677 assert_eq!(state.song_pages.lock().unwrap().clone(), vec![0]);
1678 assert_eq!(state.album_calls.lock().unwrap().len(), 3);
1679 assert_eq!(result.tracks_synced, 3);
1680 assert!(result.is_complete());
1681 }
1682
1683 #[test]
1686 fn a_failed_song_page_does_not_advance_last_sync() {
1687 let (db, _dir) = test_db();
1688 let state = Arc::new(StubState {
1689 albums: Mutex::new(stub_albums(1_200)),
1690 lists_songs: true,
1691 failing_song_pages: Mutex::new([500].into_iter().collect()),
1692 ..Default::default()
1693 });
1694 let server = StubServer::start(state.clone());
1695 let client = SubsonicClient::new(&server.url(), "u", "p");
1696
1697 let result = sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1698
1699 assert_eq!(result.pages_failed, 1);
1700 assert!(!result.is_complete());
1701 assert_eq!(get_last_sync(&db, &server.url()).unwrap(), None);
1702 let retries = state
1703 .song_pages
1704 .lock()
1705 .unwrap()
1706 .iter()
1707 .filter(|&&o| o == 500)
1708 .count();
1709 assert_eq!(retries as u32, PAGE_ATTEMPTS);
1710 }
1711
1712 #[test]
1715 fn an_incremental_sync_does_not_walk_every_song() {
1716 let (db, _dir) = test_db();
1717 let state = Arc::new(StubState {
1718 albums: Mutex::new(stub_albums(3)),
1719 lists_songs: true,
1720 ..Default::default()
1721 });
1722 let server = StubServer::start(state.clone());
1723 let client = SubsonicClient::new(&server.url(), "u", "p");
1724 update_last_sync(&db, &server.url(), "u", 0).unwrap();
1725
1726 sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1727
1728 assert!(state.song_pages.lock().unwrap().is_empty());
1729 assert_eq!(state.album_calls.lock().unwrap().len(), 3);
1730 }
1731
1732 #[test]
1733 fn sync_deduplicates_by_remote_id() {
1734 let (db, _dir) = test_db();
1735
1736 let meta1 = remote_track_meta("remote-dup", "Original Title", "Artist A", "Album X");
1738 let id1 = queries::upsert_track(&db.conn, &meta1).unwrap();
1739
1740 let meta2 = remote_track_meta("remote-dup", "Updated Title", "Artist A", "Album X");
1742 let id2 = queries::upsert_track(&db.conn, &meta2).unwrap();
1743
1744 assert_eq!(id1, id2, "same remote_id should resolve to same track row");
1746
1747 let row = queries::get_track_row(&db.conn, id2)
1749 .unwrap()
1750 .expect("track should exist");
1751 assert_eq!(row.title, "Updated Title");
1752 assert_eq!(row.remote_id.as_deref(), Some("remote-dup"));
1753
1754 let stats = queries::library_stats(&db.conn).unwrap();
1756 assert_eq!(
1757 stats.total_tracks, 1,
1758 "should have exactly 1 track after dedup"
1759 );
1760 }
1761
1762 #[test]
1766 fn resyncing_a_track_under_a_new_server_id_keeps_one_row() {
1767 let (db, _dir) = test_db();
1768 let before = remote_track_meta("46215", "Archangel", "Burial", "Untrue");
1769 let row = queries::upsert_synced_track(&db.conn, &before, &HashSet::from(["46215".into()]))
1770 .unwrap();
1771 queries::add_favourite(
1772 &db.conn,
1773 queries::LOCAL_USER,
1774 std::path::Path::new(before.remote_url.as_deref().unwrap()),
1775 )
1776 .unwrap();
1777
1778 let uid = "0199a0b2-7c4e-7d3a-9f1b-2c3d4e5f6a7b";
1779 let after = remote_track_meta(uid, "Archangel", "Burial", "Untrue");
1780 let again =
1781 queries::upsert_synced_track(&db.conn, &after, &HashSet::from([uid.into()])).unwrap();
1782
1783 assert_eq!(again, row, "the same row, not a second copy");
1784 assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 1);
1785 let (remote_id, album, artist): (String, String, String) = db
1786 .conn
1787 .query_row(
1788 "SELECT t.remote_id, al.remote_id, ar.remote_id FROM tracks t
1789 JOIN albums al ON al.id = t.album_id JOIN artists ar ON ar.id = t.artist_id",
1790 [],
1791 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1792 )
1793 .unwrap();
1794 assert_eq!(remote_id, uid);
1795 assert_eq!(album, format!("album-of-{uid}"));
1796 assert_eq!(artist, format!("artist-of-{uid}"));
1797 let adopted: String = db
1798 .conn
1799 .query_row("SELECT uid FROM tracks WHERE id = ?1", [row], |r| r.get(0))
1800 .unwrap();
1801 assert_eq!(adopted, uid, "the row takes the server's uid");
1802 let favourites = queries::load_favourites(&db.conn, queries::LOCAL_USER).unwrap();
1803 assert_eq!(
1804 favourites,
1805 HashSet::from([std::path::PathBuf::from(after.remote_url.unwrap())]),
1806 "the favourite follows the new stream address"
1807 );
1808 }
1809
1810 #[test]
1813 fn identical_entries_in_one_sync_stay_two_rows() {
1814 let (db, _dir) = test_db();
1815 let mut seen = HashSet::new();
1816 for id in ["dup-1", "dup-2"] {
1817 seen.insert(id.to_string());
1818 queries::upsert_synced_track(
1819 &db.conn,
1820 &remote_track_meta(id, "Archangel", "Burial", "Untrue"),
1821 &seen,
1822 )
1823 .unwrap();
1824 }
1825 assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 2);
1826 }
1827}