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