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