1use std::collections::HashSet;
2use std::sync::atomic::{AtomicUsize, Ordering};
3
4use crate::db::connection::Database;
5use crate::db::queries::{self, TrackMeta};
6use crate::remote::client::{SubsonicAlbumFull, SubsonicArtist, SubsonicClient};
7
8use rayon::prelude::*;
9use rusqlite::params;
10use thiserror::Error;
11
12#[derive(Debug, Error)]
13pub enum SyncError {
14 #[error("subsonic error: {0}")]
15 Subsonic(#[from] super::client::SubsonicError),
16 #[error("db error: {0}")]
17 Db(#[from] crate::db::connection::DbError),
18}
19
20#[derive(Debug, Default)]
21pub struct SyncResult {
22 pub artists_synced: usize,
23 pub albums_synced: usize,
24 pub tracks_synced: usize,
25 pub albums_failed: usize,
28}
29
30impl SyncResult {
31 pub fn is_complete(&self) -> bool {
33 self.albums_failed == 0
34 }
35}
36
37pub fn get_last_sync(
39 db: &Database,
40 url: &str,
41) -> Result<Option<i64>, crate::db::connection::DbError> {
42 let result = db.conn.query_row(
43 "SELECT last_sync FROM remote_servers WHERE url = ?1",
44 params![url],
45 |row| row.get::<_, Option<i64>>(0),
46 );
47 match result {
48 Ok(ts) => Ok(ts),
49 Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
50 Err(e) => Err(e.into()),
51 }
52}
53
54pub fn update_last_sync(
56 db: &Database,
57 url: &str,
58 username: &str,
59 timestamp: i64,
60) -> Result<(), crate::db::connection::DbError> {
61 db.conn.execute(
62 "INSERT INTO remote_servers (url, username, last_sync)
63 VALUES (?1, ?2, ?3)
64 ON CONFLICT(url) DO UPDATE SET last_sync = ?3",
65 params![url, username, timestamp],
66 )?;
67 Ok(())
68}
69
70fn parse_iso8601_to_unix(s: &str) -> Option<i64> {
78 use chrono::{DateTime, FixedOffset, NaiveDateTime};
79
80 if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
82 return Some(dt.timestamp());
83 }
84
85 if let Ok(naive) = NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S%.f") {
88 return Some(naive.and_utc().timestamp());
89 }
90 if let Ok(naive) = NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S") {
91 return Some(naive.and_utc().timestamp());
92 }
93
94 if let Ok(dt) = DateTime::<FixedOffset>::parse_from_str(s, "%Y-%m-%d %H:%M:%S%:z") {
96 return Some(dt.timestamp());
97 }
98 if let Ok(naive) = NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S") {
99 return Some(naive.and_utc().timestamp());
100 }
101
102 None
103}
104
105pub fn sync_library(
121 db: &Database,
122 client: &SubsonicClient,
123 full: bool,
124 server_url: &str,
125 username: &str,
126) -> Result<SyncResult, SyncError> {
127 let mut result = SyncResult::default();
128
129 let artists = client.get_artists()?;
130 result.artists_synced = artists.len();
131 log::info!("syncing {} artists from remote", artists.len());
132
133 let last_sync = if full {
134 None
135 } else {
136 get_last_sync(db, server_url)?
137 };
138
139 match last_sync {
140 Some(ts) => log::info!("incremental sync (albums created after {})", ts),
141 None => log::info!("full sync"),
142 }
143
144 let sync_start = std::time::SystemTime::now()
145 .duration_since(std::time::UNIX_EPOCH)
146 .unwrap_or_default()
147 .as_secs() as i64;
148
149 let mut offset = 0u32;
150 let page_size = 500u32;
151 let mut seen_ids: HashSet<String> = HashSet::new();
154
155 loop {
156 let page = client.get_album_list("alphabeticalByName", page_size, offset)?;
157 if page.is_empty() {
158 break;
159 }
160 let page_count = page.len();
161 offset += page_count as u32;
162
163 let to_fetch: Vec<String> = page
167 .into_iter()
168 .filter(|a| seen_ids.insert(a.id.clone()))
169 .filter(|a| match last_sync {
170 None => true,
171 Some(ts) => a
172 .created
173 .as_deref()
174 .and_then(parse_iso8601_to_unix)
175 .is_none_or(|created| created >= ts),
176 })
177 .map(|a| a.id)
178 .collect();
179
180 if !to_fetch.is_empty() {
181 let failures = AtomicUsize::new(0);
182 let fetched: Vec<SubsonicAlbumFull> = to_fetch
183 .into_par_iter()
184 .filter_map(|id| match client.get_album(&id) {
185 Ok(full) => Some(full),
186 Err(e) => {
187 log::warn!("failed to fetch album {}: {}", id, e);
188 failures.fetch_add(1, Ordering::Relaxed);
189 None
190 }
191 })
192 .collect();
193 result.albums_failed += failures.into_inner();
194
195 write_albums(db, client, &fetched, &mut result)?;
196 }
197
198 log::info!(
199 "synced {} albums ({} tracks) so far...",
200 result.albums_synced,
201 result.tracks_synced
202 );
203
204 if (page_count as u32) < page_size {
205 break;
206 }
207 }
208
209 write_artists(db, &artists, &mut result);
214
215 if result.is_complete() {
216 update_last_sync(db, server_url, username, sync_start)?;
217 } else {
218 log::warn!(
219 "{} album(s) failed to fetch — leaving last_sync unchanged so the next sync retries them",
220 result.albums_failed
221 );
222 }
223
224 log::info!(
225 "sync complete: {} artists, {} albums, {} tracks, {} failed",
226 result.artists_synced,
227 result.albums_synced,
228 result.tracks_synced,
229 result.albums_failed,
230 );
231
232 db.optimize();
233
234 Ok(result)
235}
236
237fn write_artists(db: &Database, artists: &[SubsonicArtist], result: &mut SyncResult) {
242 if db.conn.execute_batch("BEGIN").is_err() {
243 return;
244 }
245 let mut enriched = 0;
246 for artist in artists {
247 match queries::enrich_remote_artist(
248 &db.conn,
249 &artist.id,
250 artist.music_brainz_id.as_deref(),
251 artist.sort_name.as_deref(),
252 ) {
253 Ok(()) => enriched += 1,
254 Err(e) => log::warn!(
255 "failed to record artist metadata for {}: {}",
256 artist.name,
257 e
258 ),
259 }
260 }
261 if db.conn.execute_batch("COMMIT").is_err() {
262 let _ = db.conn.execute_batch("ROLLBACK");
263 return;
264 }
265 result.artists_synced = enriched;
266 log::info!("recorded metadata for {enriched} artists");
267}
268
269fn write_albums(
271 db: &Database,
272 client: &SubsonicClient,
273 albums: &[SubsonicAlbumFull],
274 result: &mut SyncResult,
275) -> Result<(), SyncError> {
276 db.conn
277 .execute_batch("BEGIN")
278 .map_err(crate::db::connection::DbError::from)?;
279
280 for album in albums {
281 result.albums_synced += 1;
282 let artist_name = album.artist.as_deref().unwrap_or("Unknown Artist");
283
284 for song in &album.song {
285 let meta = TrackMeta {
286 title: song.title.clone(),
287 artist: song
288 .artist
289 .clone()
290 .unwrap_or_else(|| artist_name.to_string()),
291 album_artist: album.artist.clone(),
292 album: album.name.clone(),
293 date: album.year.map(|y| y.to_string()),
294 disc: song.disc_number,
295 track_number: song.track,
296 genre: song.genre.clone().or_else(|| album.genre.clone()),
297 label: None,
298 duration_ms: song.duration.map(|d| d * 1000),
299 codec: song.suffix.clone(),
300 sample_rate: positive(song.sampling_rate),
308 bit_depth: positive(song.bit_depth),
309 channels: positive(song.channel_count),
310 bitrate: song.bit_rate,
311 size_bytes: None,
312 mtime: None,
313 path: None,
314 source: "remote".to_string(),
315 remote_id: Some(song.id.clone()),
316 remote_url: Some(client.stream_url_template(&song.id)),
317 album_remote_id: Some(album.id.clone()),
318 artist_remote_id: album.artist_id.clone(),
319 mbid: song.music_brainz_id.clone(),
320 album_mbid: album.music_brainz_id.clone(),
321 album_added_at: album.created.clone(),
322 };
323
324 match queries::upsert_track(&db.conn, &meta) {
325 Ok(_) => result.tracks_synced += 1,
326 Err(e) => log::warn!("failed to insert remote track {}: {}", song.title, e),
327 }
328 }
329
330 if let Err(e) = queries::enrich_remote_album(
334 &db.conn,
335 &album.id,
336 album.music_brainz_id.as_deref(),
337 album.sort_name.as_deref(),
338 album.song_count,
339 album.record_labels.first().map(|l| l.name.as_str()),
340 ) {
341 log::warn!("failed to record album metadata for {}: {}", album.name, e);
342 }
343 }
344
345 db.conn
346 .execute_batch("COMMIT")
347 .map_err(crate::db::connection::DbError::from)?;
348
349 Ok(())
350}
351
352fn positive(value: Option<i32>) -> Option<i32> {
355 value.filter(|v| *v > 0)
356}
357
358#[cfg(test)]
359mod tests {
360 use super::*;
361
362 #[test]
363 fn parse_rfc3339_with_z() {
364 assert_eq!(
366 parse_iso8601_to_unix("2024-01-15T10:30:00Z"),
367 Some(1705314600)
368 );
369 }
370
371 #[test]
372 fn parse_rfc3339_with_offset() {
373 assert_eq!(
375 parse_iso8601_to_unix("2024-01-15T10:30:00+05:30"),
376 Some(1705294800)
377 );
378 }
379
380 #[test]
381 fn parse_rfc3339_negative_offset() {
382 assert_eq!(
384 parse_iso8601_to_unix("2024-01-15T10:30:00-05:00"),
385 Some(1705332600)
386 );
387 }
388
389 #[test]
390 fn parse_fractional_seconds_z() {
391 assert_eq!(
392 parse_iso8601_to_unix("2024-01-15T10:30:00.123Z"),
393 Some(1705314600)
394 );
395 }
396
397 #[test]
398 fn parse_fractional_seconds_offset() {
399 assert_eq!(
400 parse_iso8601_to_unix("2024-01-15T10:30:00.999+00:00"),
401 Some(1705314600)
402 );
403 }
404
405 #[test]
406 fn parse_no_timezone_assumes_utc() {
407 assert_eq!(
408 parse_iso8601_to_unix("2024-01-15T10:30:00"),
409 Some(1705314600)
410 );
411 }
412
413 #[test]
414 fn parse_no_timezone_fractional() {
415 assert_eq!(
416 parse_iso8601_to_unix("2024-01-15T10:30:00.500"),
417 Some(1705314600)
418 );
419 }
420
421 #[test]
422 fn parse_space_separator_with_tz() {
423 assert_eq!(
424 parse_iso8601_to_unix("2024-01-15 10:30:00+00:00"),
425 Some(1705314600)
426 );
427 }
428
429 #[test]
430 fn parse_space_separator_no_tz() {
431 assert_eq!(
432 parse_iso8601_to_unix("2024-01-15 10:30:00"),
433 Some(1705314600)
434 );
435 }
436
437 #[test]
438 fn parse_garbage_returns_none() {
439 assert_eq!(parse_iso8601_to_unix("not-a-date"), None);
440 assert_eq!(parse_iso8601_to_unix(""), None);
441 assert_eq!(parse_iso8601_to_unix("2024"), None);
442 }
443
444 #[test]
445 fn parse_epoch() {
446 assert_eq!(parse_iso8601_to_unix("1970-01-01T00:00:00Z"), Some(0));
447 }
448
449 use crate::db::connection::Database;
452 use crate::db::queries;
453 use std::sync::{Arc, Mutex};
454
455 fn test_db() -> (Database, tempfile::TempDir) {
456 let dir = tempfile::tempdir().unwrap();
457 let db = Database::open(&dir.path().join("sync_test.db")).unwrap();
458 (db, dir)
459 }
460
461 fn remote_track_meta(remote_id: &str, title: &str, artist: &str, album: &str) -> TrackMeta {
463 TrackMeta {
464 title: title.into(),
465 artist: artist.into(),
466 album_artist: Some(artist.into()),
467 album: album.into(),
468 date: Some("2024".into()),
469 disc: Some(1),
470 track_number: Some(1),
471 genre: Some("Electronic".into()),
472 label: None,
473 duration_ms: Some(240_000),
474 codec: Some("FLAC".into()),
475 sample_rate: Some(44100),
476 bit_depth: Some(16),
477 channels: Some(2),
478 bitrate: Some(1000),
479 size_bytes: None,
480 mtime: None,
481 path: None,
482 source: "remote".into(),
483 remote_id: Some(remote_id.into()),
484 remote_url: Some(format!("https://example.com/stream?id={}", remote_id)),
485 album_remote_id: Some(format!("album-of-{remote_id}")),
486 artist_remote_id: Some(format!("artist-of-{remote_id}")),
487 mbid: Some(format!("mbid-of-{remote_id}")),
488 album_mbid: None,
489 album_added_at: None,
490 }
491 }
492
493 #[test]
496 fn a_zero_quality_figure_is_treated_as_absent() {
497 assert_eq!(positive(Some(0)), None);
498 assert_eq!(positive(Some(16)), Some(16));
499 assert_eq!(positive(None), None);
500 }
501
502 #[test]
506 fn album_metadata_from_the_server_is_recorded() {
507 let (db, _dir) = test_db();
508
509 let meta = remote_track_meta("remote-300", "Anguish", "Sleep", "Volume One");
510 queries::upsert_track(&db.conn, &meta).unwrap();
511
512 queries::enrich_remote_album(
513 &db.conn,
514 "album-of-remote-300",
515 Some("mb-album-1"),
516 Some("volume one"),
517 Some(6),
518 Some("Off The Disk"),
519 )
520 .unwrap();
521
522 let row: (Option<String>, Option<String>, Option<i32>, Option<String>) = db
523 .conn
524 .query_row(
525 "SELECT mbid, sort_name, total_tracks, label FROM albums WHERE title = 'Volume One'",
526 [],
527 |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
528 )
529 .unwrap();
530 assert_eq!(row.0.as_deref(), Some("mb-album-1"));
531 assert_eq!(row.1.as_deref(), Some("volume one"));
532 assert_eq!(row.2, Some(6));
533 assert_eq!(row.3.as_deref(), Some("Off The Disk"));
534 }
535
536 #[test]
539 fn server_metadata_does_not_overwrite_what_tags_said() {
540 let (db, _dir) = test_db();
541
542 let mut meta = remote_track_meta("remote-301", "Dopesmoker", "Sleep", "Dopesmoker");
543 meta.label = Some("From The Tags".into());
544 queries::upsert_track(&db.conn, &meta).unwrap();
545
546 queries::enrich_remote_album(
547 &db.conn,
548 "album-of-remote-301",
549 None,
550 None,
551 None,
552 Some("From The Server"),
553 )
554 .unwrap();
555
556 let label: Option<String> = db
557 .conn
558 .query_row(
559 "SELECT label FROM albums WHERE title = 'Dopesmoker'",
560 [],
561 |r| r.get(0),
562 )
563 .unwrap();
564 assert_eq!(label.as_deref(), Some("From The Tags"));
565 }
566
567 #[test]
570 fn artist_metadata_from_the_server_is_recorded() {
571 let (db, _dir) = test_db();
572
573 let meta = remote_track_meta("remote-302", "Holy Mountain", "Sleep", "Holy Mountain");
574 queries::upsert_track(&db.conn, &meta).unwrap();
575
576 queries::enrich_remote_artist(
577 &db.conn,
578 "artist-of-remote-302",
579 Some("mb-artist-1"),
580 Some("sleep"),
581 )
582 .unwrap();
583
584 let row: (Option<String>, Option<String>) = db
585 .conn
586 .query_row(
587 "SELECT mbid, sort_name FROM artists WHERE name = 'Sleep'",
588 [],
589 |r| Ok((r.get(0)?, r.get(1)?)),
590 )
591 .unwrap();
592 assert_eq!(row.0.as_deref(), Some("mb-artist-1"));
593 assert_eq!(row.1.as_deref(), Some("sleep"));
594 }
595
596 #[test]
598 fn a_synced_track_keeps_its_musicbrainz_id() {
599 let (db, _dir) = test_db();
600
601 let meta = remote_track_meta("remote-303", "Aquarian", "Sleep", "Dopesmoker");
602 let id = queries::upsert_track(&db.conn, &meta).unwrap();
603
604 let mbid: Option<String> = db
605 .conn
606 .query_row("SELECT mbid FROM tracks WHERE id = ?1", [id], |r| r.get(0))
607 .unwrap();
608 assert_eq!(mbid.as_deref(), Some("mbid-of-remote-303"));
609 }
610
611 #[test]
615 fn a_synced_track_keeps_its_quality_figures() {
616 let (db, _dir) = test_db();
617
618 let meta = remote_track_meta("remote-200", "Anguish", "Sleep", "Volume One");
619 let id = queries::upsert_track(&db.conn, &meta).unwrap();
620
621 let row = queries::get_track_row(&db.conn, id).unwrap().unwrap();
622 assert_eq!(row.sample_rate, Some(44100));
623 assert_eq!(row.bit_depth, Some(16));
624 assert_eq!(row.channels, Some(2));
625 }
626
627 #[test]
631 fn a_sync_records_the_album_and_artist_ids_too() {
632 let (db, _dir) = test_db();
633
634 let meta = remote_track_meta("remote-100", "Enter", "Russian Circles", "Enter");
635 queries::upsert_track(&db.conn, &meta).unwrap();
636
637 let album: Option<String> = db
638 .conn
639 .query_row(
640 "SELECT remote_id FROM albums WHERE title = 'Enter'",
641 [],
642 |r| r.get(0),
643 )
644 .unwrap();
645 assert_eq!(album.as_deref(), Some("album-of-remote-100"));
646
647 let artist: Option<String> = db
648 .conn
649 .query_row(
650 "SELECT remote_id FROM artists WHERE name = 'Russian Circles'",
651 [],
652 |r| r.get(0),
653 )
654 .unwrap();
655 assert_eq!(artist.as_deref(), Some("artist-of-remote-100"));
656 }
657
658 #[test]
659 fn sync_upserts_tracks_to_database() {
660 let (db, _dir) = test_db();
661
662 let meta = remote_track_meta("remote-001", "Vordhosbn", "Aphex Twin", "Drukqs");
663 let track_id = queries::upsert_track(&db.conn, &meta).unwrap();
664 assert!(track_id > 0, "upsert should return a valid track ID");
665
666 let row = queries::get_track_row(&db.conn, track_id)
668 .unwrap()
669 .expect("track should exist in DB");
670 assert_eq!(row.title, "Vordhosbn");
671 assert_eq!(row.artist_name, "Aphex Twin");
672 assert_eq!(row.album_title, "Drukqs");
673 assert_eq!(row.remote_id.as_deref(), Some("remote-001"));
674 assert_eq!(row.source, "remote");
675 }
676
677 struct StubServer {
682 addr: std::net::SocketAddr,
683 shutdown: Arc<std::sync::atomic::AtomicBool>,
684 }
685
686 #[derive(Default)]
687 struct StubState {
688 albums: Mutex<Vec<(String, String, String)>>,
690 failing: Mutex<HashSet<String>>,
692 insert_after_first_page: Mutex<Option<(String, String, String)>>,
695 list_pages_served: AtomicUsize,
696 list_types: Mutex<Vec<String>>,
697 album_calls: Mutex<Vec<String>>,
698 }
699
700 impl StubServer {
701 fn start(state: Arc<StubState>) -> Self {
702 let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
703 listener.set_nonblocking(true).unwrap();
704 let addr = listener.local_addr().unwrap();
705 let shutdown = Arc::new(std::sync::atomic::AtomicBool::new(false));
706
707 let stop = shutdown.clone();
708 std::thread::spawn(move || {
709 while !stop.load(Ordering::Relaxed) {
710 match listener.accept() {
711 Ok((stream, _)) => {
712 let _ = stream.set_nonblocking(false);
714 let state = state.clone();
715 std::thread::spawn(move || handle(stream, state));
716 }
717 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
718 std::thread::sleep(std::time::Duration::from_millis(2));
719 }
720 Err(_) => break,
721 }
722 }
723 });
724
725 Self { addr, shutdown }
726 }
727
728 fn url(&self) -> String {
729 format!("http://{}", self.addr)
730 }
731 }
732
733 impl Drop for StubServer {
734 fn drop(&mut self) {
735 self.shutdown.store(true, Ordering::Relaxed);
736 }
737 }
738
739 fn handle(mut stream: std::net::TcpStream, state: Arc<StubState>) {
743 use std::io::{BufRead, Write};
744
745 let Ok(peek) = stream.try_clone() else { return };
746 let mut reader = std::io::BufReader::new(peek);
747
748 loop {
749 let mut request_line = String::new();
750 if reader.read_line(&mut request_line).unwrap_or(0) == 0 {
751 return;
752 }
753 let mut line = String::new();
754 while reader.read_line(&mut line).unwrap_or(0) > 0 {
755 if line == "\r\n" || line == "\n" {
756 break;
757 }
758 line.clear();
759 }
760
761 let target = request_line.split_whitespace().nth(1).unwrap_or("/");
762 let (path, query) = target.split_once('?').unwrap_or((target, ""));
763 let params: std::collections::HashMap<&str, &str> = query
764 .split('&')
765 .filter_map(|kv| kv.split_once('='))
766 .collect();
767
768 let (status, body) = respond(&state, path, ¶ms);
769
770 let write = write!(
771 stream,
772 "HTTP/1.1 {} OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n",
773 status,
774 body.len()
775 )
776 .and_then(|()| stream.write_all(body.as_bytes()))
777 .and_then(|()| stream.flush());
778 if write.is_err() {
779 return;
780 }
781 }
782 }
783
784 fn respond(
785 state: &Arc<StubState>,
786 path: &str,
787 params: &std::collections::HashMap<&str, &str>,
788 ) -> (u16, String) {
789 match path.rsplit('/').next().unwrap_or("") {
790 "getArtists" => (
791 200,
792 r#"{"subsonic-response":{"status":"ok","artists":{"index":[{"artist":[{"id":"ar1","name":"Stub Artist"}]}]}}}"#.to_string(),
793 ),
794 "getAlbumList2" => {
795 state
796 .list_types
797 .lock()
798 .unwrap()
799 .push(params.get("type").copied().unwrap_or("").to_string());
800 let offset: usize = params.get("offset").and_then(|o| o.parse().ok()).unwrap_or(0);
801 let size: usize = params.get("size").and_then(|s| s.parse().ok()).unwrap_or(500);
802
803 let albums = state.albums.lock().unwrap();
804 let slice: Vec<String> = albums
805 .iter()
806 .skip(offset)
807 .take(size)
808 .map(|(id, name, created)| {
809 format!(
810 r#"{{"id":"{}","name":"{}","artist":"Stub Artist","created":"{}"}}"#,
811 id, name, created
812 )
813 })
814 .collect();
815 drop(albums);
816
817 if state.list_pages_served.fetch_add(1, Ordering::SeqCst) == 0
818 && let Some(new_album) = state.insert_after_first_page.lock().unwrap().take()
819 {
820 state.albums.lock().unwrap().insert(0, new_album);
821 }
822
823 (
824 200,
825 format!(
826 r#"{{"subsonic-response":{{"status":"ok","albumList2":{{"album":[{}]}}}}}}"#,
827 slice.join(",")
828 ),
829 )
830 }
831 "getAlbum" => {
832 let id = params.get("id").copied().unwrap_or("");
833 state.album_calls.lock().unwrap().push(id.to_string());
834 if state.failing.lock().unwrap().contains(id) {
835 (500, r#"{"error":"boom"}"#.to_string())
836 } else {
837 (
838 200,
839 format!(
840 r#"{{"subsonic-response":{{"status":"ok","album":{{"id":"{id}","name":"Album {id}","artist":"Stub Artist","song":[{{"id":"s{id}","title":"Song {id}","track":1,"suffix":"flac"}}]}}}}}}"#
841 ),
842 )
843 }
844 }
845 _ => (404, "{}".to_string()),
846 }
847 }
848
849 fn stub_albums(n: usize) -> Vec<(String, String, String)> {
850 (0..n)
851 .map(|i| {
852 (
853 format!("a{:04}", i),
854 format!("Album {:04}", i),
855 "2024-01-15T10:30:00Z".to_string(),
856 )
857 })
858 .collect()
859 }
860
861 #[test]
862 fn failed_album_fetch_does_not_advance_last_sync_and_next_sync_retries() {
863 let (db, _dir) = test_db();
864 let state = Arc::new(StubState {
865 albums: Mutex::new(stub_albums(4)),
866 failing: Mutex::new(["a0002".to_string()].into_iter().collect()),
867 ..Default::default()
868 });
869 let server = StubServer::start(state.clone());
870 let client = SubsonicClient::new(&server.url(), "u", "p");
871
872 let first = sync_library(&db, &client, false, &server.url(), "u").unwrap();
873 assert_eq!(first.albums_failed, 1, "the failing album must be counted");
874 assert_eq!(first.albums_synced, 3);
875 assert!(!first.is_complete());
876 assert_eq!(
877 get_last_sync(&db, &server.url()).unwrap(),
878 None,
879 "an incomplete sync must not advance last_sync"
880 );
881
882 state.failing.lock().unwrap().clear();
885 state.album_calls.lock().unwrap().clear();
886 let second = sync_library(&db, &client, false, &server.url(), "u").unwrap();
887
888 assert!(
889 state
890 .album_calls
891 .lock()
892 .unwrap()
893 .contains(&"a0002".to_string()),
894 "the previously failed album must be retried"
895 );
896 assert_eq!(second.albums_failed, 0);
897 assert!(second.is_complete());
898 assert!(
899 get_last_sync(&db, &server.url()).unwrap().is_some(),
900 "a clean sync advances last_sync"
901 );
902 }
903
904 #[test]
905 fn album_inserted_mid_pagination_is_not_fetched_twice_or_skipped() {
906 let (db, _dir) = test_db();
910 let state = Arc::new(StubState {
911 albums: Mutex::new(stub_albums(600)),
912 insert_after_first_page: Mutex::new(Some((
913 "aNEW".to_string(),
914 "AAA Brand New".to_string(),
915 "2024-06-01T00:00:00Z".to_string(),
916 ))),
917 ..Default::default()
918 });
919 let server = StubServer::start(state.clone());
920 let client = SubsonicClient::new(&server.url(), "u", "p");
921
922 let result = sync_library(&db, &client, true, &server.url(), "u").unwrap();
923 assert_eq!(result.albums_failed, 0);
924
925 let calls = state.album_calls.lock().unwrap().clone();
926 let unique: HashSet<&String> = calls.iter().collect();
927 assert_eq!(
928 calls.len(),
929 unique.len(),
930 "no album may be fetched twice after the window shifts"
931 );
932
933 for i in 0..600 {
935 let id = format!("a{:04}", i);
936 assert!(unique.contains(&id), "album {} was skipped", id);
937 }
938
939 let types = state.list_types.lock().unwrap().clone();
940 assert!(
941 types.iter().all(|t| t == "alphabeticalByName"),
942 "the paginated walk must use a stable ordering, got {:?}",
943 types
944 );
945 }
946
947 #[test]
948 fn incremental_sync_only_fetches_albums_created_after_last_sync() {
949 let (db, _dir) = test_db();
950 let mut albums = stub_albums(3);
951 albums[0].2 = "2020-01-01T00:00:00Z".into();
952 albums[1].2 = "2020-01-01T00:00:00Z".into();
953 albums[2].2 = "2030-01-01T00:00:00Z".into();
954
955 let state = Arc::new(StubState {
956 albums: Mutex::new(albums),
957 ..Default::default()
958 });
959 let server = StubServer::start(state.clone());
960 let client = SubsonicClient::new(&server.url(), "u", "p");
961
962 let watermark = parse_iso8601_to_unix("2025-01-01T00:00:00Z").unwrap();
964 update_last_sync(&db, &server.url(), "u", watermark).unwrap();
965
966 let result = sync_library(&db, &client, false, &server.url(), "u").unwrap();
967
968 assert_eq!(result.albums_synced, 1, "only the new album needs fetching");
969 assert_eq!(
970 *state.album_calls.lock().unwrap(),
971 vec!["a0002".to_string()]
972 );
973 }
974
975 #[test]
976 fn album_with_unparseable_created_is_always_fetched() {
977 let (db, _dir) = test_db();
978 let state = Arc::new(StubState {
979 albums: Mutex::new(vec![("a0000".into(), "Album".into(), "who knows".into())]),
980 ..Default::default()
981 });
982 let server = StubServer::start(state.clone());
983 let client = SubsonicClient::new(&server.url(), "u", "p");
984
985 update_last_sync(&db, &server.url(), "u", 4_000_000_000).unwrap();
986 let result = sync_library(&db, &client, false, &server.url(), "u").unwrap();
987
988 assert_eq!(
989 result.albums_synced, 1,
990 "an album with no usable timestamp must not be assumed old"
991 );
992 }
993
994 #[test]
995 fn sync_deduplicates_by_remote_id() {
996 let (db, _dir) = test_db();
997
998 let meta1 = remote_track_meta("remote-dup", "Original Title", "Artist A", "Album X");
1000 let id1 = queries::upsert_track(&db.conn, &meta1).unwrap();
1001
1002 let meta2 = remote_track_meta("remote-dup", "Updated Title", "Artist A", "Album X");
1004 let id2 = queries::upsert_track(&db.conn, &meta2).unwrap();
1005
1006 assert_eq!(id1, id2, "same remote_id should resolve to same track row");
1008
1009 let row = queries::get_track_row(&db.conn, id2)
1011 .unwrap()
1012 .expect("track should exist");
1013 assert_eq!(row.title, "Updated Title");
1014 assert_eq!(row.remote_id.as_deref(), Some("remote-dup"));
1015
1016 let stats = queries::library_stats(&db.conn).unwrap();
1018 assert_eq!(
1019 stats.total_tracks, 1,
1020 "should have exactly 1 track after dedup"
1021 );
1022 }
1023}