Skip to main content

koan_core/remote/
sync.rs

1use std::collections::{HashMap, HashSet};
2use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
3
4use crate::db::connection::Database;
5use crate::db::queries::{self, TrackMeta};
6use crate::remote::client::{
7    SubsonicAlbum, SubsonicAlbumFull, SubsonicArtist, SubsonicClient, SubsonicSong,
8};
9
10use rayon::prelude::*;
11use rusqlite::params;
12use thiserror::Error;
13
14#[derive(Debug, Error)]
15pub enum SyncError {
16    #[error("subsonic error: {0}")]
17    Subsonic(#[from] super::client::SubsonicError),
18    #[error("db error: {0}")]
19    Db(#[from] crate::db::connection::DbError),
20}
21
22#[derive(Debug, Default)]
23pub struct SyncResult {
24    pub artists_synced: usize,
25    pub albums_synced: usize,
26    pub tracks_synced: usize,
27    /// Albums whose details could not be fetched. Non-zero means `last_sync`
28    /// was left where it was so the next sync picks them up again.
29    pub albums_failed: usize,
30    /// Pages of songs that could not be fetched, on the bulk path. Counts
31    /// against completeness the same way.
32    pub pages_failed: usize,
33    /// Tracks removed because the server no longer has them.
34    pub tracks_removed: usize,
35}
36
37impl SyncResult {
38    /// Whether the run covered everything it set out to.
39    pub fn is_complete(&self) -> bool {
40        self.albums_failed == 0 && self.pages_failed == 0
41    }
42}
43
44/// Get the last sync timestamp for a remote server, if any.
45pub 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
61/// Update (or insert) the last sync timestamp for a remote server.
62pub 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
77/// Parse an ISO 8601 / RFC 3339 timestamp string into a unix timestamp (seconds).
78/// Returns `None` if the string can't be parsed.
79///
80/// Handles common Subsonic/Navidrome variants:
81/// - Full RFC 3339: `2024-01-15T10:30:00Z`, `2024-01-15T10:30:00+05:30`
82/// - Fractional seconds: `2024-01-15T10:30:00.123Z`
83/// - Missing timezone (assumed UTC): `2024-01-15T10:30:00`
84fn parse_iso8601_to_unix(s: &str) -> Option<i64> {
85    use chrono::{DateTime, FixedOffset, NaiveDateTime};
86
87    // Try strict RFC 3339 first (handles Z, offsets, fractional seconds).
88    if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
89        return Some(dt.timestamp());
90    }
91
92    // Subsonic sometimes omits timezone — parse as naive and assume UTC.
93    // Try with fractional seconds first, then without.
94    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    // Some servers use space instead of T.
102    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/// Which part of a sync is running.
113#[derive(Debug, Clone, Copy, PartialEq, Eq)]
114pub enum SyncPhase {
115    /// Paging through the server's album list.
116    Albums,
117    /// Fetching and writing tracks. The long part.
118    Tracks,
119    /// Recording artist metadata.
120    Artists,
121    /// Relinking, recording the watermark, optimising the database.
122    Finishing,
123}
124
125/// How far a sync has got. `done` and `total` count albums in the `Albums`
126/// phase and tracks in the `Tracks` phase; `total` is `None` where the server
127/// gives no way to know it in advance.
128#[derive(Debug, Clone, Copy, PartialEq, Eq)]
129pub struct SyncProgress {
130    pub phase: SyncPhase,
131    pub done: u64,
132    pub total: Option<u64>,
133}
134
135/// Albums and songs per list request. 500 is the most `getAlbumList2` allows.
136const PAGE_SIZE: u32 = 500;
137
138/// Song pages in flight at once during a full sync. The walk is bound by
139/// round trips, not bandwidth; four hides most of a mobile link's latency
140/// without asking much of the server.
141const FETCH_LANES: usize = 4;
142
143/// Attempts at one song page before it counts as failed.
144const PAGE_ATTEMPTS: u32 = 3;
145
146/// Changed albums above which an incremental sync pages every song rather
147/// than fetching each album: about a hundred requests for the whole library,
148/// against one per album.
149const PAGED_ABOVE: usize = 200;
150
151/// Pull the Navidrome/Subsonic library into the local DB.
152///
153/// The album list is always walked in `alphabeticalByName` order: it is the one
154/// ordering stable under concurrent server-side inserts, so an offset walk can
155/// never skip an album that was added between two pages.
156///
157/// A first or full sync then pages through every song with an empty `search3`
158/// query — about a hundred requests for fifty thousand tracks — and joins them
159/// to the album list. Fetching each album on its own took one round trip per
160/// album, which on a phone was minutes. A server that does not answer an empty
161/// query gets the per-album walk instead, as does an incremental sync, which
162/// only fetches albums created after `last_sync` and so has few to fetch.
163///
164/// `last_sync` only advances when everything was fetched. A run that lost
165/// albums or pages to network errors leaves the timestamp alone so the next
166/// sync fetches them again, rather than writing a permanent hole in the library.
167///
168/// Deduplication happens in `upsert_track` — if a local track already exists
169/// with the same artist + album + title + track#, the remote_id and remote_url
170/// are merged onto the existing row instead of creating a duplicate.
171pub fn sync_library(
172    db: &Database,
173    client: &SubsonicClient,
174    full: bool,
175    server_url: &str,
176    username: &str,
177    progress: &(dyn Fn(SyncProgress) + Sync),
178) -> Result<SyncResult, SyncError> {
179    let mut result = SyncResult::default();
180
181    let last_sync = if full {
182        None
183    } else {
184        get_last_sync(db, server_url)?
185    };
186
187    match last_sync {
188        Some(ts) => log::info!("incremental sync (albums created after {})", ts),
189        None => log::info!("full sync"),
190    }
191
192    let sync_start = std::time::SystemTime::now()
193        .duration_since(std::time::UNIX_EPOCH)
194        .unwrap_or_default()
195        .as_secs() as i64;
196
197    let albums = list_albums(client, progress)?;
198
199    // An incremental sync fetches albums created since the last one, and any
200    // the client holds differently from how the server lists them: a retag
201    // on the server keeps an album's `created` and can give it a new id, and
202    // neither would otherwise ever be read again. An unparseable `created` is
203    // treated as new: re-fetching is cheap, missing is not.
204    let held = last_sync.and_then(|_| held_albums(db));
205    let wanted: Vec<&SubsonicAlbum> = albums
206        .iter()
207        .filter(|a| match last_sync {
208            None => true,
209            Some(ts) => {
210                a.created
211                    .as_deref()
212                    .and_then(parse_iso8601_to_unix)
213                    .is_none_or(|created| created >= ts)
214                    || held.as_ref().is_some_and(|h| differs(a, h.get(&a.id)))
215            }
216        })
217        .collect();
218    let expected: u64 = wanted
219        .iter()
220        .filter_map(|a| a.song_count)
221        .map(|n| n.max(0) as u64)
222        .sum();
223
224    let mut song_ids: HashSet<String> = HashSet::new();
225    let mut total = (expected > 0).then_some(expected);
226    let walked = if (last_sync.is_none() || wanted.len() > PAGED_ABOVE) && !wanted.is_empty() {
227        if total.is_none() {
228            total = client.song_count().ok().flatten();
229        }
230        sync_all_songs(
231            db,
232            client,
233            &albums,
234            total,
235            &mut result,
236            &mut song_ids,
237            progress,
238        )?
239    } else {
240        false
241    };
242    if !walked {
243        sync_by_album(
244            db,
245            client,
246            &wanted,
247            total,
248            &mut result,
249            &mut song_ids,
250            progress,
251        )?;
252    }
253
254    // Artists only ever existed as a side effect of a track upsert, which is
255    // why not one of them had a MusicBrainz id or a sort name. Applied last,
256    // because the rows do not exist until their tracks have been written.
257    let artists = client.get_artists()?;
258    result.artists_synced = artists.len();
259    progress(SyncProgress {
260        phase: SyncPhase::Artists,
261        done: 0,
262        total: Some(artists.len() as u64),
263    });
264    write_artists(db, &artists, &mut result);
265
266    progress(SyncProgress {
267        phase: SyncPhase::Finishing,
268        done: 0,
269        total: None,
270    });
271
272    // An empty listing is far likelier to be a server fault than an empty
273    // library, and would unlink every file. So is one shorter than the count
274    // the server gave: a track deleted mid-walk shifts a later page by one, and
275    // the track it pushes out of view is not gone.
276    let listed_everything = total.is_none_or(|n| song_ids.len() as u64 >= n);
277    if full && result.is_complete() && !song_ids.is_empty() && listed_everything {
278        match queries::relink_vanished_remote_ids(&db.conn, &song_ids) {
279            Ok(0) => {}
280            Ok(n) => log::info!("{n} files had ids the server no longer knows; relinked"),
281            Err(e) => log::warn!("failed to relink tracks with vanished remote ids: {e}"),
282        }
283    }
284
285    // What the server deleted goes here too. Every sync lists every album, so
286    // an album gone from that list goes on any sync; a single track deleted
287    // from an album only shows on a full one, which lists every track. Both
288    // are held to the same guards as above: an empty or short listing is a
289    // fault, not a deletion.
290    let live_albums: std::collections::HashSet<String> =
291        albums.iter().map(|a| a.id.clone()).collect();
292    let live_tracks = (full && result.is_complete() && !song_ids.is_empty() && listed_everything)
293        .then_some(&song_ids);
294    if result.is_complete() && !live_albums.is_empty() {
295        match queries::remove_vanished_remote(&db.conn, live_tracks, Some(&live_albums)) {
296            Ok(0) => {}
297            Ok(n) => {
298                result.tracks_removed = n;
299                log::info!("{n} tracks the server no longer has; removed");
300            }
301            Err(e) => log::warn!("failed to remove tracks the server deleted: {e}"),
302        }
303    }
304
305    if result.is_complete() {
306        update_last_sync(db, server_url, username, sync_start)?;
307    } else {
308        log::warn!(
309            "{} album(s) and {} page(s) failed to fetch — leaving last_sync unchanged so the next sync retries them",
310            result.albums_failed,
311            result.pages_failed,
312        );
313    }
314
315    log::info!(
316        "sync complete: {} artists, {} albums, {} tracks, {} failed",
317        result.artists_synced,
318        result.albums_synced,
319        result.tracks_synced,
320        result.albums_failed,
321    );
322
323    db.optimize();
324
325    Ok(result)
326}
327
328/// An album as this client holds it: what a listing of the server's albums
329/// can be checked against.
330struct HeldAlbum {
331    title: String,
332    artist: String,
333    tracks: i64,
334    seconds: i64,
335}
336
337/// Every album this client holds from the server, by the server's id. `None`
338/// if it cannot be read, which leaves an incremental sync to go by `created`
339/// alone rather than fetching everything.
340fn held_albums(db: &Database) -> Option<HashMap<String, HeldAlbum>> {
341    let read = || -> rusqlite::Result<HashMap<String, HeldAlbum>> {
342        let mut stmt = db.conn.prepare(
343            "SELECT al.remote_id, al.title, COALESCE(ar.name, ''),
344                    COUNT(t.id), COALESCE(SUM(t.duration_ms), 0) / 1000
345             FROM albums al
346             LEFT JOIN artists ar ON ar.id = al.artist_id
347             LEFT JOIN tracks t ON t.album_id = al.id AND t.remote_id IS NOT NULL
348             WHERE al.remote_id IS NOT NULL
349             GROUP BY al.id",
350        )?;
351        stmt.query_map([], |r| {
352            Ok((
353                r.get::<_, String>(0)?,
354                HeldAlbum {
355                    title: r.get(1)?,
356                    artist: r.get(2)?,
357                    tracks: r.get(3)?,
358                    seconds: r.get(4)?,
359                },
360            ))
361        })?
362        .collect()
363    };
364    read()
365        .inspect_err(|e| log::warn!("could not read held albums; syncing by date alone: {e}"))
366        .ok()
367}
368
369/// Whether the server lists an album differently from how it is held: not
370/// held at all, renamed, credited to someone else, or with other tracks.
371fn differs(listed: &SubsonicAlbum, held: Option<&HeldAlbum>) -> bool {
372    let Some(held) = held else { return true };
373    listed.name != held.title
374        || listed.artist.as_deref().is_some_and(|a| a != held.artist)
375        || listed
376            .song_count
377            .is_some_and(|n| i64::from(n) != held.tracks)
378        || listed
379            .duration
380            .is_some_and(|d| (d - held.seconds).abs() > 2)
381}
382
383/// Every album on the server, once each.
384fn list_albums(
385    client: &SubsonicClient,
386    progress: &(dyn Fn(SyncProgress) + Sync),
387) -> Result<Vec<SubsonicAlbum>, SyncError> {
388    let mut albums = Vec::new();
389    // Guards against an album appearing on two pages when the server-side list
390    // shifts under the offset walk.
391    let mut seen: HashSet<String> = HashSet::new();
392    let mut offset = 0u32;
393    loop {
394        let page = client.get_album_list("alphabeticalByName", PAGE_SIZE, offset)?;
395        let count = page.len() as u32;
396        offset += count;
397        albums.extend(page.into_iter().filter(|a| seen.insert(a.id.clone())));
398        progress(SyncProgress {
399            phase: SyncPhase::Albums,
400            done: albums.len() as u64,
401            total: None,
402        });
403        if count < PAGE_SIZE {
404            return Ok(albums);
405        }
406    }
407}
408
409/// Page through every song with an empty `search3` query, joined to the album
410/// list, one transaction per page.
411///
412/// Returns `false`, having written nothing, when the server does not list
413/// songs that way — older servers answer an empty query with nothing, or with
414/// an error. The caller then walks the albums one at a time.
415fn sync_all_songs(
416    db: &Database,
417    client: &SubsonicClient,
418    albums: &[SubsonicAlbum],
419    total: Option<u64>,
420    result: &mut SyncResult,
421    song_ids: &mut HashSet<String>,
422    progress: &(dyn Fn(SyncProgress) + Sync),
423) -> Result<bool, SyncError> {
424    let first = match client.all_songs_page(PAGE_SIZE, 0) {
425        Ok(songs) if !songs.is_empty() => songs,
426        Ok(_) => {
427            log::info!("server lists no songs for an empty search; syncing album by album");
428            return Ok(false);
429        }
430        Err(e) => {
431            log::info!("server refused an empty search ({e}); syncing album by album");
432            return Ok(false);
433        }
434    };
435
436    let by_id: HashMap<&str, &SubsonicAlbum> = albums.iter().map(|a| (a.id.as_str(), a)).collect();
437    let mut albums_seen: HashSet<String> = HashSet::new();
438    let mut write = |songs: Vec<SubsonicSong>,
439                     result: &mut SyncResult,
440                     song_ids: &mut HashSet<String>|
441     -> Result<(), SyncError> {
442        let batch = group_by_album(songs, &by_id);
443        albums_seen.extend(batch.iter().map(|a| a.id.clone()));
444        write_albums(db, client, &batch, result, song_ids)?;
445        progress(SyncProgress {
446            phase: SyncPhase::Tracks,
447            done: song_ids.len() as u64,
448            total,
449        });
450        Ok(())
451    };
452
453    let short = (first.len() as u32) < PAGE_SIZE;
454    write(first, result, song_ids)?;
455
456    if !short {
457        // Pages are fetched on a few lanes and written here, in whatever order
458        // they arrive — each is a whole transaction on its own, and the
459        // database has one writer however many are fetching. The lanes are
460        // threads rather than rayon tasks because they spend their time
461        // waiting on the network, not the CPU.
462        let next = AtomicU32::new(PAGE_SIZE);
463        // The offset of the first page that came back short. Nothing past it
464        // is worth asking for.
465        let end = AtomicU32::new(u32::MAX);
466        let (tx, rx) = std::sync::mpsc::sync_channel(FETCH_LANES);
467        let failed = std::thread::scope(|scope| -> Result<usize, SyncError> {
468            for _ in 0..FETCH_LANES {
469                let tx = tx.clone();
470                let (next, end) = (&next, &end);
471                scope.spawn(move || {
472                    loop {
473                        let offset = next.fetch_add(PAGE_SIZE, Ordering::Relaxed);
474                        if offset >= end.load(Ordering::Relaxed) {
475                            return;
476                        }
477                        let page = fetch_page(client, offset);
478                        // A page that failed after its retries ends the walk
479                        // too: the server is likely gone, and asking on for
480                        // ever would never finish. The run is then incomplete
481                        // and the next sync walks it again.
482                        if page
483                            .as_ref()
484                            .map_or(true, |songs| (songs.len() as u32) < PAGE_SIZE)
485                        {
486                            end.fetch_min(offset, Ordering::Relaxed);
487                        }
488                        if tx.send((offset, page)).is_err() {
489                            return;
490                        }
491                    }
492                });
493            }
494            drop(tx);
495
496            let mut failed = 0;
497            for (offset, page) in rx {
498                match page {
499                    Ok(songs) => write(songs, result, song_ids)?,
500                    Err(e) => {
501                        log::warn!("failed to fetch songs from offset {offset}: {e}");
502                        failed += 1;
503                    }
504                }
505            }
506            Ok(failed)
507        })?;
508        result.pages_failed += failed;
509    }
510
511    result.albums_synced += albums_seen.len();
512    Ok(true)
513}
514
515/// One page of songs, retried a couple of times: a failed page is a hole in
516/// the library until the next full sync, so it is worth a second ask.
517fn fetch_page(
518    client: &SubsonicClient,
519    offset: u32,
520) -> Result<Vec<SubsonicSong>, super::client::SubsonicError> {
521    let mut attempt = 1;
522    loop {
523        match client.all_songs_page(PAGE_SIZE, offset) {
524            Ok(songs) => return Ok(songs),
525            Err(e) if attempt >= PAGE_ATTEMPTS => return Err(e),
526            Err(e) => {
527                log::debug!("songs from offset {offset}, attempt {attempt}: {e}");
528                std::thread::sleep(std::time::Duration::from_millis(250 * attempt as u64));
529                attempt += 1;
530            }
531        }
532    }
533}
534
535/// A page of songs as the albums they belong to, each carrying the metadata
536/// the album list gave for it.
537///
538/// A song whose album is not in the list — added after the list was read — is
539/// written under what the song itself says about its album.
540fn group_by_album(
541    songs: Vec<SubsonicSong>,
542    albums: &HashMap<&str, &SubsonicAlbum>,
543) -> Vec<SubsonicAlbumFull> {
544    let mut grouped: Vec<SubsonicAlbumFull> = Vec::new();
545    let mut index: HashMap<String, usize> = HashMap::new();
546    for song in songs {
547        let Some(album_id) = song.album_id.clone() else {
548            log::warn!("song {} has no album id; skipped", song.id);
549            continue;
550        };
551        let i = *index.entry(album_id.clone()).or_insert_with(|| {
552            grouped.push(match albums.get(album_id.as_str()) {
553                Some(album) => SubsonicAlbumFull {
554                    id: album.id.clone(),
555                    name: album.name.clone(),
556                    artist: album.artist.clone(),
557                    artist_id: album.artist_id.clone(),
558                    year: album.year,
559                    genre: album.genre.clone(),
560                    song_count: album.song_count,
561                    created: album.created.clone(),
562                    music_brainz_id: album.music_brainz_id.clone(),
563                    sort_name: album.sort_name.clone(),
564                    record_labels: album.record_labels.clone(),
565                    song: Vec::new(),
566                },
567                None => SubsonicAlbumFull {
568                    id: album_id,
569                    name: song.album.clone().unwrap_or_default(),
570                    artist: song.artist.clone(),
571                    artist_id: song.artist_id.clone(),
572                    year: song.year,
573                    genre: song.genre.clone(),
574                    song_count: None,
575                    created: None,
576                    music_brainz_id: None,
577                    sort_name: None,
578                    record_labels: Vec::new(),
579                    song: Vec::new(),
580                },
581            });
582            grouped.len() - 1
583        });
584        grouped[i].song.push(song);
585    }
586    grouped
587}
588
589/// Albums to write per transaction on the per-album path. Small enough that
590/// progress moves often, large enough that the fetches overlap.
591const ALBUM_BATCH: usize = 100;
592
593/// Fetch each album on its own and write them a batch at a time. What an
594/// incremental sync does, and what a full one falls back to on a server that
595/// cannot list songs in bulk.
596fn sync_by_album(
597    db: &Database,
598    client: &SubsonicClient,
599    albums: &[&SubsonicAlbum],
600    total: Option<u64>,
601    result: &mut SyncResult,
602    song_ids: &mut HashSet<String>,
603    progress: &(dyn Fn(SyncProgress) + Sync),
604) -> Result<(), SyncError> {
605    for batch in albums.chunks(ALBUM_BATCH) {
606        let failures = AtomicUsize::new(0);
607        let fetched: Vec<SubsonicAlbumFull> = batch
608            .par_iter()
609            .filter_map(|album| match client.get_album(&album.id) {
610                Ok(full) => Some(full),
611                Err(e) => {
612                    log::warn!("failed to fetch album {}: {}", album.id, e);
613                    failures.fetch_add(1, Ordering::Relaxed);
614                    None
615                }
616            })
617            .collect();
618        result.albums_failed += failures.into_inner();
619        result.albums_synced += fetched.len();
620
621        write_albums(db, client, &fetched, result, song_ids)?;
622        progress(SyncProgress {
623            phase: SyncPhase::Tracks,
624            done: song_ids.len() as u64,
625            total,
626        });
627        log::info!(
628            "synced {} albums ({} tracks) so far...",
629            result.albums_synced,
630            result.tracks_synced
631        );
632    }
633    Ok(())
634}
635
636/// Record what the server knows about each artist.
637///
638/// Best-effort: a library that synced its tracks fine should not fail because
639/// one artist row could not be updated.
640fn write_artists(db: &Database, artists: &[SubsonicArtist], result: &mut SyncResult) {
641    if db.conn.execute_batch("BEGIN").is_err() {
642        return;
643    }
644    let mut enriched = 0;
645    for artist in artists {
646        match queries::enrich_remote_artist(
647            &db.conn,
648            &artist.id,
649            artist.music_brainz_id.as_deref(),
650            artist.sort_name.as_deref(),
651        ) {
652            Ok(()) => enriched += 1,
653            Err(e) => log::warn!(
654                "failed to record artist metadata for {}: {}",
655                artist.name,
656                e
657            ),
658        }
659    }
660    if db.conn.execute_batch("COMMIT").is_err() {
661        let _ = db.conn.execute_batch("ROLLBACK");
662        return;
663    }
664    result.artists_synced = enriched;
665    log::info!("recorded metadata for {enriched} artists");
666}
667
668/// Write one batch of fetched albums in a single transaction.
669fn write_albums(
670    db: &Database,
671    client: &SubsonicClient,
672    albums: &[SubsonicAlbumFull],
673    result: &mut SyncResult,
674    song_ids: &mut HashSet<String>,
675) -> Result<(), SyncError> {
676    db.conn
677        .execute_batch("BEGIN")
678        .map_err(crate::db::connection::DbError::from)?;
679
680    for album in albums {
681        let artist_name = album.artist.as_deref().unwrap_or("Unknown Artist");
682
683        for song in &album.song {
684            song_ids.insert(song.id.clone());
685            let meta = TrackMeta {
686                title: song.title.clone(),
687                artist: song
688                    .artist
689                    .clone()
690                    .unwrap_or_else(|| artist_name.to_string()),
691                album_artist: album.artist.clone(),
692                album: album.name.clone(),
693                date: album.year.map(|y| y.to_string()),
694                disc: song.disc_number,
695                track_number: song.track,
696                genre: song.genre.clone().or_else(|| album.genre.clone()),
697                label: None,
698                duration_ms: song.duration.map(|d| d * 1000),
699                codec: song.suffix.clone(),
700                // OpenSubsonic servers report these; a plain Subsonic one
701                // leaves them out and the track keeps no quality figures,
702                // which is what every remote track used to get.
703                //
704                // Zero means "not applicable", not "zero" — Navidrome reports
705                // bitDepth 0 for every lossy file. Storing it would render an
706                // MP3 as 0-bit, which is worse than saying nothing.
707                sample_rate: positive(song.sampling_rate),
708                bit_depth: positive(song.bit_depth),
709                channels: positive(song.channel_count),
710                bitrate: song.bit_rate,
711                size_bytes: None,
712                mtime: None,
713                path: None,
714                source: "remote".to_string(),
715                remote_id: Some(song.id.clone()),
716                remote_url: Some(client.stream_url_template(&song.id)),
717                album_remote_id: Some(album.id.clone()),
718                artist_remote_id: album.artist_id.clone(),
719                mbid: song.music_brainz_id.clone(),
720                album_mbid: album.music_brainz_id.clone(),
721                album_added_at: album.created.clone(),
722            };
723
724            match queries::upsert_synced_track(&db.conn, &meta, song_ids) {
725                Ok(_) => result.tracks_synced += 1,
726                Err(e) => log::warn!("failed to insert remote track {}: {}", song.title, e),
727            }
728        }
729
730        // The album row is created through its tracks, so it only ever sees
731        // what a file's tags say. Track totals, the label and the MusicBrainz
732        // id belong to the release, and came back in the same response.
733        if let Err(e) = queries::enrich_remote_album(
734            &db.conn,
735            &album.id,
736            album.music_brainz_id.as_deref(),
737            album.sort_name.as_deref(),
738            album.song_count,
739            album
740                .record_labels
741                .iter()
742                .map(|l| l.name.as_str())
743                .find(|l| !l.is_empty()),
744        ) {
745            log::warn!("failed to record album metadata for {}: {}", album.name, e);
746        }
747    }
748
749    db.conn
750        .execute_batch("COMMIT")
751        .map_err(crate::db::connection::DbError::from)?;
752
753    Ok(())
754}
755
756/// Treat a non-positive figure as absent. None of sample rate, bit depth or
757/// channel count has a meaningful zero.
758fn positive(value: Option<i32>) -> Option<i32> {
759    value.filter(|v| *v > 0)
760}
761
762#[cfg(test)]
763mod tests {
764    use super::*;
765
766    #[test]
767    fn parse_rfc3339_with_z() {
768        // 2024-01-15T10:30:00Z = 1705314600
769        assert_eq!(
770            parse_iso8601_to_unix("2024-01-15T10:30:00Z"),
771            Some(1705314600)
772        );
773    }
774
775    #[test]
776    fn parse_rfc3339_with_offset() {
777        // 10:30 IST (+05:30) = 05:00 UTC = 1705294800
778        assert_eq!(
779            parse_iso8601_to_unix("2024-01-15T10:30:00+05:30"),
780            Some(1705294800)
781        );
782    }
783
784    #[test]
785    fn parse_rfc3339_negative_offset() {
786        // 10:30 EST (-05:00) = 15:30 UTC = 1705332600
787        assert_eq!(
788            parse_iso8601_to_unix("2024-01-15T10:30:00-05:00"),
789            Some(1705332600)
790        );
791    }
792
793    #[test]
794    fn parse_fractional_seconds_z() {
795        assert_eq!(
796            parse_iso8601_to_unix("2024-01-15T10:30:00.123Z"),
797            Some(1705314600)
798        );
799    }
800
801    #[test]
802    fn parse_fractional_seconds_offset() {
803        assert_eq!(
804            parse_iso8601_to_unix("2024-01-15T10:30:00.999+00:00"),
805            Some(1705314600)
806        );
807    }
808
809    #[test]
810    fn parse_no_timezone_assumes_utc() {
811        assert_eq!(
812            parse_iso8601_to_unix("2024-01-15T10:30:00"),
813            Some(1705314600)
814        );
815    }
816
817    #[test]
818    fn parse_no_timezone_fractional() {
819        assert_eq!(
820            parse_iso8601_to_unix("2024-01-15T10:30:00.500"),
821            Some(1705314600)
822        );
823    }
824
825    #[test]
826    fn parse_space_separator_with_tz() {
827        assert_eq!(
828            parse_iso8601_to_unix("2024-01-15 10:30:00+00:00"),
829            Some(1705314600)
830        );
831    }
832
833    #[test]
834    fn parse_space_separator_no_tz() {
835        assert_eq!(
836            parse_iso8601_to_unix("2024-01-15 10:30:00"),
837            Some(1705314600)
838        );
839    }
840
841    #[test]
842    fn parse_garbage_returns_none() {
843        assert_eq!(parse_iso8601_to_unix("not-a-date"), None);
844        assert_eq!(parse_iso8601_to_unix(""), None);
845        assert_eq!(parse_iso8601_to_unix("2024"), None);
846    }
847
848    #[test]
849    fn parse_epoch() {
850        assert_eq!(parse_iso8601_to_unix("1970-01-01T00:00:00Z"), Some(0));
851    }
852
853    // --- Sync → DB integration tests ---
854
855    use crate::db::connection::Database;
856    use crate::db::queries;
857    use std::sync::{Arc, Mutex};
858
859    fn test_db() -> (Database, tempfile::TempDir) {
860        let dir = tempfile::tempdir().unwrap();
861        let db = Database::open(&dir.path().join("sync_test.db")).unwrap();
862        (db, dir)
863    }
864
865    /// Build a TrackMeta matching how sync_library constructs them from SubsonicSong data.
866    fn remote_track_meta(remote_id: &str, title: &str, artist: &str, album: &str) -> TrackMeta {
867        TrackMeta {
868            title: title.into(),
869            artist: artist.into(),
870            album_artist: Some(artist.into()),
871            album: album.into(),
872            date: Some("2024".into()),
873            disc: Some(1),
874            track_number: Some(1),
875            genre: Some("Electronic".into()),
876            label: None,
877            duration_ms: Some(240_000),
878            codec: Some("FLAC".into()),
879            sample_rate: Some(44100),
880            bit_depth: Some(16),
881            channels: Some(2),
882            bitrate: Some(1000),
883            size_bytes: None,
884            mtime: None,
885            path: None,
886            source: "remote".into(),
887            remote_id: Some(remote_id.into()),
888            remote_url: Some(format!("https://example.com/stream?id={}", remote_id)),
889            album_remote_id: Some(format!("album-of-{remote_id}")),
890            artist_remote_id: Some(format!("artist-of-{remote_id}")),
891            mbid: Some(format!("mbid-of-{remote_id}")),
892            album_mbid: None,
893            album_added_at: None,
894        }
895    }
896
897    /// Navidrome reports bitDepth 0 for every lossy file. Keeping it would
898    /// render an MP3 as 0-bit; absent is the honest answer.
899    #[test]
900    fn a_zero_quality_figure_is_treated_as_absent() {
901        assert_eq!(positive(Some(0)), None);
902        assert_eq!(positive(Some(16)), Some(16));
903        assert_eq!(positive(None), None);
904    }
905
906    /// The album row is created through its tracks, so it only ever sees what
907    /// a file's tags say. Track totals, the label and the MusicBrainz id are
908    /// properties of the release and arrive in the same response.
909    #[test]
910    fn album_metadata_from_the_server_is_recorded() {
911        let (db, _dir) = test_db();
912
913        let meta = remote_track_meta("remote-300", "Anguish", "Sleep", "Volume One");
914        queries::upsert_track(&db.conn, &meta).unwrap();
915
916        queries::enrich_remote_album(
917            &db.conn,
918            "album-of-remote-300",
919            Some("mb-album-1"),
920            Some("volume one"),
921            Some(6),
922            Some("Off The Disk"),
923        )
924        .unwrap();
925
926        let row: (Option<String>, Option<String>, Option<i32>, Option<String>) = db
927            .conn
928            .query_row(
929                "SELECT mbid, sort_name, total_tracks, label FROM albums WHERE title = 'Volume One'",
930                [],
931                |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
932            )
933            .unwrap();
934        assert_eq!(row.0.as_deref(), Some("mb-album-1"));
935        assert_eq!(row.1.as_deref(), Some("volume one"));
936        assert_eq!(row.2, Some(6));
937        assert_eq!(row.3.as_deref(), Some("Off The Disk"));
938    }
939
940    /// Enrichment fills blanks. A locally-scanned album whose tags named a
941    /// label must keep it, rather than have the server's answer written over.
942    #[test]
943    fn server_metadata_does_not_overwrite_what_tags_said() {
944        let (db, _dir) = test_db();
945
946        let mut meta = remote_track_meta("remote-301", "Dopesmoker", "Sleep", "Dopesmoker");
947        meta.label = Some("From The Tags".into());
948        queries::upsert_track(&db.conn, &meta).unwrap();
949
950        queries::enrich_remote_album(
951            &db.conn,
952            "album-of-remote-301",
953            None,
954            None,
955            None,
956            Some("From The Server"),
957        )
958        .unwrap();
959
960        let label: Option<String> = db
961            .conn
962            .query_row(
963                "SELECT label FROM albums WHERE title = 'Dopesmoker'",
964                [],
965                |r| r.get(0),
966            )
967            .unwrap();
968        assert_eq!(label.as_deref(), Some("From The Tags"));
969    }
970
971    /// Artists existed only as a side effect of a track upsert, so nothing
972    /// ever wrote their MusicBrainz id or sort name.
973    #[test]
974    fn artist_metadata_from_the_server_is_recorded() {
975        let (db, _dir) = test_db();
976
977        let meta = remote_track_meta("remote-302", "Holy Mountain", "Sleep", "Holy Mountain");
978        queries::upsert_track(&db.conn, &meta).unwrap();
979
980        queries::enrich_remote_artist(
981            &db.conn,
982            "artist-of-remote-302",
983            Some("mb-artist-1"),
984            Some("sleep"),
985        )
986        .unwrap();
987
988        let row: (Option<String>, Option<String>) = db
989            .conn
990            .query_row(
991                "SELECT mbid, sort_name FROM artists WHERE name = 'Sleep'",
992                [],
993                |r| Ok((r.get(0)?, r.get(1)?)),
994            )
995            .unwrap();
996        assert_eq!(row.0.as_deref(), Some("mb-artist-1"));
997        assert_eq!(row.1.as_deref(), Some("sleep"));
998    }
999
1000    /// The recording id travels with the track.
1001    #[test]
1002    fn a_synced_track_keeps_its_musicbrainz_id() {
1003        let (db, _dir) = test_db();
1004
1005        let meta = remote_track_meta("remote-303", "Aquarian", "Sleep", "Dopesmoker");
1006        let id = queries::upsert_track(&db.conn, &meta).unwrap();
1007
1008        let mbid: Option<String> = db
1009            .conn
1010            .query_row("SELECT mbid FROM tracks WHERE id = ?1", [id], |r| r.get(0))
1011            .unwrap();
1012        assert_eq!(mbid.as_deref(), Some("mbid-of-remote-303"));
1013    }
1014
1015    /// A remote track carries the quality figures an OpenSubsonic server
1016    /// reports. Without them the format badge has nothing to show, which is
1017    /// what every remote-only track in a synced library used to look like.
1018    #[test]
1019    fn a_synced_track_keeps_its_quality_figures() {
1020        let (db, _dir) = test_db();
1021
1022        let meta = remote_track_meta("remote-200", "Anguish", "Sleep", "Volume One");
1023        let id = queries::upsert_track(&db.conn, &meta).unwrap();
1024
1025        let row = queries::get_track_row(&db.conn, id).unwrap().unwrap();
1026        assert_eq!(row.sample_rate, Some(44100));
1027        assert_eq!(row.bit_depth, Some(16));
1028        assert_eq!(row.channels, Some(2));
1029    }
1030
1031    /// The server keys stars, shares and cover art off album and artist ids,
1032    /// so a sync that only kept the track's left the library unable to refer
1033    /// to either — every album row came back with a null remote_id.
1034    #[test]
1035    fn a_sync_records_the_album_and_artist_ids_too() {
1036        let (db, _dir) = test_db();
1037
1038        let meta = remote_track_meta("remote-100", "Enter", "Russian Circles", "Enter");
1039        queries::upsert_track(&db.conn, &meta).unwrap();
1040
1041        let album: Option<String> = db
1042            .conn
1043            .query_row(
1044                "SELECT remote_id FROM albums WHERE title = 'Enter'",
1045                [],
1046                |r| r.get(0),
1047            )
1048            .unwrap();
1049        assert_eq!(album.as_deref(), Some("album-of-remote-100"));
1050
1051        let artist: Option<String> = db
1052            .conn
1053            .query_row(
1054                "SELECT remote_id FROM artists WHERE name = 'Russian Circles'",
1055                [],
1056                |r| r.get(0),
1057            )
1058            .unwrap();
1059        assert_eq!(artist.as_deref(), Some("artist-of-remote-100"));
1060    }
1061
1062    #[test]
1063    fn sync_upserts_tracks_to_database() {
1064        let (db, _dir) = test_db();
1065
1066        let meta = remote_track_meta("remote-001", "Vordhosbn", "Aphex Twin", "Drukqs");
1067        let track_id = queries::upsert_track(&db.conn, &meta).unwrap();
1068        assert!(track_id > 0, "upsert should return a valid track ID");
1069
1070        // Verify the track exists with correct remote_id.
1071        let row = queries::get_track_row(&db.conn, track_id)
1072            .unwrap()
1073            .expect("track should exist in DB");
1074        assert_eq!(row.title, "Vordhosbn");
1075        assert_eq!(row.artist_name, "Aphex Twin");
1076        assert_eq!(row.album_title, "Drukqs");
1077        assert_eq!(row.remote_id.as_deref(), Some("remote-001"));
1078        assert_eq!(row.source, "remote");
1079    }
1080
1081    // --- sync_library against a stub Subsonic server ---
1082
1083    /// Minimal Subsonic server: serves getArtists, getAlbumList2 and getAlbum
1084    /// so `sync_library` can be driven end to end without a real Navidrome.
1085    struct StubServer {
1086        addr: std::net::SocketAddr,
1087        shutdown: Arc<std::sync::atomic::AtomicBool>,
1088    }
1089
1090    #[derive(Default)]
1091    struct StubState {
1092        /// Album (id, name, created) in the order the server lists them.
1093        albums: Mutex<Vec<(String, String, String)>>,
1094        /// Album ids whose getAlbum call fails with a 500.
1095        failing: Mutex<HashSet<String>>,
1096        /// Prepended to the album list once the first list page has been served,
1097        /// modelling a server-side insert landing mid-pagination.
1098        insert_after_first_page: Mutex<Option<(String, String, String)>>,
1099        list_pages_served: AtomicUsize,
1100        list_types: Mutex<Vec<String>>,
1101        album_calls: Mutex<Vec<String>>,
1102        /// Answer an empty `search3` with every album's song, as an
1103        /// OpenSubsonic server does. Off, it answers with nothing, as older
1104        /// servers do.
1105        lists_songs: bool,
1106        /// `songOffset`s asked for.
1107        song_pages: Mutex<Vec<usize>>,
1108        /// Song pages that fail with a 500, by offset.
1109        failing_song_pages: Mutex<HashSet<usize>>,
1110    }
1111
1112    impl StubServer {
1113        fn start(state: Arc<StubState>) -> Self {
1114            let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
1115            listener.set_nonblocking(true).unwrap();
1116            let addr = listener.local_addr().unwrap();
1117            let shutdown = Arc::new(std::sync::atomic::AtomicBool::new(false));
1118
1119            let stop = shutdown.clone();
1120            std::thread::spawn(move || {
1121                while !stop.load(Ordering::Relaxed) {
1122                    match listener.accept() {
1123                        Ok((stream, _)) => {
1124                            // BSD sockets inherit O_NONBLOCK from the listener.
1125                            let _ = stream.set_nonblocking(false);
1126                            let state = state.clone();
1127                            std::thread::spawn(move || handle(stream, state));
1128                        }
1129                        Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
1130                            std::thread::sleep(std::time::Duration::from_millis(2));
1131                        }
1132                        Err(_) => break,
1133                    }
1134                }
1135            });
1136
1137            Self { addr, shutdown }
1138        }
1139
1140        fn url(&self) -> String {
1141            format!("http://{}", self.addr)
1142        }
1143    }
1144
1145    impl Drop for StubServer {
1146        fn drop(&mut self) {
1147            self.shutdown.store(true, Ordering::Relaxed);
1148        }
1149    }
1150
1151    /// Serve requests on one connection until the peer closes it. Keep-alive
1152    /// matters here: reqwest pools connections, and a server that hangs up after
1153    /// every response makes parallel fetches fail on reused sockets.
1154    fn handle(mut stream: std::net::TcpStream, state: Arc<StubState>) {
1155        use std::io::{BufRead, Write};
1156
1157        let Ok(peek) = stream.try_clone() else { return };
1158        let mut reader = std::io::BufReader::new(peek);
1159
1160        loop {
1161            let mut request_line = String::new();
1162            if reader.read_line(&mut request_line).unwrap_or(0) == 0 {
1163                return;
1164            }
1165            let mut line = String::new();
1166            while reader.read_line(&mut line).unwrap_or(0) > 0 {
1167                if line == "\r\n" || line == "\n" {
1168                    break;
1169                }
1170                line.clear();
1171            }
1172
1173            let target = request_line.split_whitespace().nth(1).unwrap_or("/");
1174            let (path, query) = target.split_once('?').unwrap_or((target, ""));
1175            let params: std::collections::HashMap<&str, &str> = query
1176                .split('&')
1177                .filter_map(|kv| kv.split_once('='))
1178                .collect();
1179
1180            let (status, body) = respond(&state, path, &params);
1181
1182            let write = write!(
1183                stream,
1184                "HTTP/1.1 {} OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n",
1185                status,
1186                body.len()
1187            )
1188            .and_then(|()| stream.write_all(body.as_bytes()))
1189            .and_then(|()| stream.flush());
1190            if write.is_err() {
1191                return;
1192            }
1193        }
1194    }
1195
1196    fn respond(
1197        state: &Arc<StubState>,
1198        path: &str,
1199        params: &std::collections::HashMap<&str, &str>,
1200    ) -> (u16, String) {
1201        match path.rsplit('/').next().unwrap_or("") {
1202            "getArtists" => (
1203                200,
1204                r#"{"subsonic-response":{"status":"ok","artists":{"index":[{"artist":[{"id":"ar1","name":"Stub Artist"}]}]}}}"#.to_string(),
1205            ),
1206            "getAlbumList2" => {
1207                state
1208                    .list_types
1209                    .lock()
1210                    .unwrap()
1211                    .push(params.get("type").copied().unwrap_or("").to_string());
1212                let offset: usize = params.get("offset").and_then(|o| o.parse().ok()).unwrap_or(0);
1213                let size: usize = params.get("size").and_then(|s| s.parse().ok()).unwrap_or(500);
1214
1215                let albums = state.albums.lock().unwrap();
1216                let slice: Vec<String> = albums
1217                    .iter()
1218                    .skip(offset)
1219                    .take(size)
1220                    .map(|(id, name, created)| {
1221                        format!(
1222                            r#"{{"id":"{}","name":"{}","artist":"Stub Artist","created":"{}","songCount":1}}"#,
1223                            id, name, created
1224                        )
1225                    })
1226                    .collect();
1227                drop(albums);
1228
1229                if state.list_pages_served.fetch_add(1, Ordering::SeqCst) == 0
1230                    && let Some(new_album) = state.insert_after_first_page.lock().unwrap().take()
1231                {
1232                    state.albums.lock().unwrap().insert(0, new_album);
1233                }
1234
1235                (
1236                    200,
1237                    format!(
1238                        r#"{{"subsonic-response":{{"status":"ok","albumList2":{{"album":[{}]}}}}}}"#,
1239                        slice.join(",")
1240                    ),
1241                )
1242            }
1243            "search3" => {
1244                let offset: usize = params.get("songOffset").and_then(|o| o.parse().ok()).unwrap_or(0);
1245                let size: usize = params.get("songCount").and_then(|s| s.parse().ok()).unwrap_or(20);
1246                state.song_pages.lock().unwrap().push(offset);
1247                if state.failing_song_pages.lock().unwrap().contains(&offset) {
1248                    return (500, r#"{"error":"boom"}"#.to_string());
1249                }
1250                let songs: Vec<String> = if state.lists_songs {
1251                    state
1252                        .albums
1253                        .lock()
1254                        .unwrap()
1255                        .iter()
1256                        .skip(offset)
1257                        .take(size)
1258                        .map(|(id, name, _)| {
1259                            format!(
1260                                r#"{{"id":"s{id}","title":"Song {id}","albumId":"{id}","album":"{name}","track":1,"suffix":"flac"}}"#
1261                            )
1262                        })
1263                        .collect()
1264                } else {
1265                    Vec::new()
1266                };
1267                (
1268                    200,
1269                    format!(
1270                        r#"{{"subsonic-response":{{"status":"ok","searchResult3":{{"song":[{}]}}}}}}"#,
1271                        songs.join(",")
1272                    ),
1273                )
1274            }
1275            "getAlbum" => {
1276                let id = params.get("id").copied().unwrap_or("");
1277                state.album_calls.lock().unwrap().push(id.to_string());
1278                if state.failing.lock().unwrap().contains(id) {
1279                    (500, r#"{"error":"boom"}"#.to_string())
1280                } else {
1281                    (
1282                        200,
1283                        format!(
1284                            r#"{{"subsonic-response":{{"status":"ok","album":{{"id":"{id}","name":"{name}","artist":"Stub Artist","song":[{{"id":"s{id}","title":"Song {id}","track":1,"suffix":"flac"}}]}}}}}}"#,
1285                            name = state
1286                                .albums
1287                                .lock()
1288                                .unwrap()
1289                                .iter()
1290                                .find(|(a, _, _)| a == id)
1291                                .map_or_else(|| format!("Album {id}"), |(_, n, _)| n.clone())
1292                        ),
1293                    )
1294                }
1295            }
1296            _ => (404, "{}".to_string()),
1297        }
1298    }
1299
1300    fn stub_albums(n: usize) -> Vec<(String, String, String)> {
1301        (0..n)
1302            .map(|i| {
1303                (
1304                    format!("a{:04}", i),
1305                    format!("Album {:04}", i),
1306                    "2024-01-15T10:30:00Z".to_string(),
1307                )
1308            })
1309            .collect()
1310    }
1311
1312    #[test]
1313    fn failed_album_fetch_does_not_advance_last_sync_and_next_sync_retries() {
1314        let (db, _dir) = test_db();
1315        let state = Arc::new(StubState {
1316            albums: Mutex::new(stub_albums(4)),
1317            failing: Mutex::new(["a0002".to_string()].into_iter().collect()),
1318            ..Default::default()
1319        });
1320        let server = StubServer::start(state.clone());
1321        let client = SubsonicClient::new(&server.url(), "u", "p");
1322
1323        let first = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1324        assert_eq!(first.albums_failed, 1, "the failing album must be counted");
1325        assert_eq!(first.albums_synced, 3);
1326        assert!(!first.is_complete());
1327        assert_eq!(
1328            get_last_sync(&db, &server.url()).unwrap(),
1329            None,
1330            "an incomplete sync must not advance last_sync"
1331        );
1332
1333        // Second run: the album now succeeds and is picked up because the sync
1334        // still has no watermark to skip past.
1335        state.failing.lock().unwrap().clear();
1336        state.album_calls.lock().unwrap().clear();
1337        let second = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1338
1339        assert!(
1340            state
1341                .album_calls
1342                .lock()
1343                .unwrap()
1344                .contains(&"a0002".to_string()),
1345            "the previously failed album must be retried"
1346        );
1347        assert_eq!(second.albums_failed, 0);
1348        assert!(second.is_complete());
1349        assert!(
1350            get_last_sync(&db, &server.url()).unwrap().is_some(),
1351            "a clean sync advances last_sync"
1352        );
1353    }
1354
1355    #[test]
1356    fn album_inserted_mid_pagination_is_not_fetched_twice_or_skipped() {
1357        // 600 albums forces a second list page; a server-side insert between
1358        // pages shifts the offset window, which without de-dup replays the
1359        // page boundary and, with `newest` ordering, can drop albums entirely.
1360        let (db, _dir) = test_db();
1361        let state = Arc::new(StubState {
1362            albums: Mutex::new(stub_albums(600)),
1363            insert_after_first_page: Mutex::new(Some((
1364                "aNEW".to_string(),
1365                "AAA Brand New".to_string(),
1366                "2024-06-01T00:00:00Z".to_string(),
1367            ))),
1368            ..Default::default()
1369        });
1370        let server = StubServer::start(state.clone());
1371        let client = SubsonicClient::new(&server.url(), "u", "p");
1372
1373        let result = sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1374        assert_eq!(result.albums_failed, 0);
1375
1376        let calls = state.album_calls.lock().unwrap().clone();
1377        let unique: HashSet<&String> = calls.iter().collect();
1378        assert_eq!(
1379            calls.len(),
1380            unique.len(),
1381            "no album may be fetched twice after the window shifts"
1382        );
1383
1384        // Every album present before the shift must still have been fetched.
1385        for i in 0..600 {
1386            let id = format!("a{:04}", i);
1387            assert!(unique.contains(&id), "album {} was skipped", id);
1388        }
1389
1390        let types = state.list_types.lock().unwrap().clone();
1391        assert!(
1392            types.iter().all(|t| t == "alphabeticalByName"),
1393            "the paginated walk must use a stable ordering, got {:?}",
1394            types
1395        );
1396    }
1397
1398    #[test]
1399    fn a_full_sync_relinks_files_whose_server_id_changed() {
1400        let (db, _dir) = test_db();
1401        let file = TrackMeta {
1402            date: None,
1403            disc: None,
1404            path: Some("/music/song.flac".into()),
1405            source: "local".into(),
1406            album_remote_id: None,
1407            artist_remote_id: None,
1408            mbid: None,
1409            ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1410        };
1411        queries::upsert_track(&db.conn, &file).unwrap();
1412
1413        let state = Arc::new(StubState {
1414            albums: Mutex::new(stub_albums(1)),
1415            ..Default::default()
1416        });
1417        let server = StubServer::start(state);
1418        let client = SubsonicClient::new(&server.url(), "u", "p");
1419        sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1420
1421        let rows: Vec<(Option<String>, Option<String>)> = db
1422            .conn
1423            .prepare("SELECT path, remote_id FROM tracks WHERE title = 'Song a0000'")
1424            .unwrap()
1425            .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1426            .unwrap()
1427            .collect::<Result<_, _>>()
1428            .unwrap();
1429        assert_eq!(
1430            rows,
1431            vec![(Some("/music/song.flac".into()), Some("sa0000".into()))],
1432            "the file takes the server's current id, and the recording is listed once"
1433        );
1434    }
1435
1436    #[test]
1437    fn a_full_sync_folds_server_rows_whose_id_vanished() {
1438        let (db, _dir) = test_db();
1439        let ghost = TrackMeta {
1440            date: None,
1441            disc: None,
1442            album_remote_id: None,
1443            artist_remote_id: None,
1444            mbid: None,
1445            ..remote_track_meta("s-before-rescan", "Song a0000", "Stub Artist", "Album 0000")
1446        };
1447        let ghost_id = queries::upsert_track(&db.conn, &ghost).unwrap();
1448        let ghost_url = ghost.remote_url.clone().unwrap();
1449        db.conn
1450            .execute(
1451                "INSERT INTO favourites (track_path) VALUES (?1)",
1452                params![ghost_url],
1453            )
1454            .unwrap();
1455
1456        let state = Arc::new(StubState {
1457            albums: Mutex::new(stub_albums(1)),
1458            ..Default::default()
1459        });
1460        let server = StubServer::start(state);
1461        let client = SubsonicClient::new(&server.url(), "u", "p");
1462        sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1463
1464        let rows: Vec<(i64, String, String)> = db
1465            .conn
1466            .prepare("SELECT id, remote_id, remote_url FROM tracks WHERE title = 'Song a0000'")
1467            .unwrap()
1468            .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))
1469            .unwrap()
1470            .collect::<Result<_, _>>()
1471            .unwrap();
1472        assert_eq!(rows.len(), 1, "the dead row takes the live id");
1473        let (id, remote_id, remote_url) = &rows[0];
1474        assert_eq!(*id, ghost_id);
1475        assert_eq!(remote_id, "sa0000");
1476        let favourites: Vec<String> = db
1477            .conn
1478            .prepare("SELECT track_path FROM favourites")
1479            .unwrap()
1480            .query_map([], |row| row.get(0))
1481            .unwrap()
1482            .collect::<Result<_, _>>()
1483            .unwrap();
1484        assert_eq!(
1485            favourites,
1486            vec![remote_url.clone()],
1487            "the favourite follows"
1488        );
1489    }
1490
1491    #[test]
1492    fn incremental_sync_only_fetches_albums_created_after_last_sync() {
1493        let (db, _dir) = test_db();
1494        let mut albums = stub_albums(3);
1495        albums[0].2 = "2020-01-01T00:00:00Z".into();
1496        albums[1].2 = "2020-01-01T00:00:00Z".into();
1497        albums[2].2 = "2030-01-01T00:00:00Z".into();
1498
1499        let state = Arc::new(StubState {
1500            albums: Mutex::new(albums),
1501            ..Default::default()
1502        });
1503        let server = StubServer::start(state.clone());
1504        let client = SubsonicClient::new(&server.url(), "u", "p");
1505
1506        // Everything held as the server lists it, then the watermark put
1507        // between the two vintages.
1508        sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1509        state.album_calls.lock().unwrap().clear();
1510        let watermark = parse_iso8601_to_unix("2025-01-01T00:00:00Z").unwrap();
1511        update_last_sync(&db, &server.url(), "u", watermark).unwrap();
1512
1513        let result = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1514
1515        assert_eq!(result.albums_synced, 1, "only the new album needs fetching");
1516        assert_eq!(
1517            *state.album_calls.lock().unwrap(),
1518            vec!["a0002".to_string()]
1519        );
1520    }
1521
1522    /// A retag on the server keeps an album's `created`. Going by the date
1523    /// alone, the client would hold the old tags until a full sync.
1524    #[test]
1525    fn incremental_sync_fetches_an_album_the_server_now_lists_differently() {
1526        let (db, _dir) = test_db();
1527        let mut albums = stub_albums(3);
1528        for a in &mut albums {
1529            a.2 = "2020-01-01T00:00:00Z".into();
1530        }
1531        let state = Arc::new(StubState {
1532            albums: Mutex::new(albums),
1533            ..Default::default()
1534        });
1535        let server = StubServer::start(state.clone());
1536        let client = SubsonicClient::new(&server.url(), "u", "p");
1537        sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1538        state.album_calls.lock().unwrap().clear();
1539
1540        state.albums.lock().unwrap()[1].1 = "Album 0001 (Retitled)".into();
1541        let result = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1542
1543        assert_eq!(
1544            *state.album_calls.lock().unwrap(),
1545            vec!["a0001".to_string()]
1546        );
1547        assert_eq!(result.albums_synced, 1);
1548        let title: String = db
1549            .conn
1550            .query_row(
1551                "SELECT title FROM albums WHERE remote_id = 'a0001'",
1552                [],
1553                |r| r.get(0),
1554            )
1555            .unwrap();
1556        assert_eq!(title, "Album 0001 (Retitled)");
1557    }
1558
1559    /// An album the client has never held is fetched however old the server
1560    /// says it is: a retag can move tracks to a new album id with an old date.
1561    #[test]
1562    fn incremental_sync_fetches_an_album_it_does_not_hold() {
1563        let (db, _dir) = test_db();
1564        let mut albums = stub_albums(2);
1565        for a in &mut albums {
1566            a.2 = "2020-01-01T00:00:00Z".into();
1567        }
1568        let state = Arc::new(StubState {
1569            albums: Mutex::new(albums),
1570            ..Default::default()
1571        });
1572        let server = StubServer::start(state.clone());
1573        let client = SubsonicClient::new(&server.url(), "u", "p");
1574        sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1575        state.album_calls.lock().unwrap().clear();
1576
1577        state.albums.lock().unwrap().push((
1578            "a0099".into(),
1579            "Album 0099".into(),
1580            "2020-01-01T00:00:00Z".into(),
1581        ));
1582        sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1583
1584        assert_eq!(
1585            *state.album_calls.lock().unwrap(),
1586            vec!["a0099".to_string()]
1587        );
1588    }
1589
1590    #[test]
1591    fn album_with_unparseable_created_is_always_fetched() {
1592        let (db, _dir) = test_db();
1593        let state = Arc::new(StubState {
1594            albums: Mutex::new(vec![("a0000".into(), "Album".into(), "who knows".into())]),
1595            ..Default::default()
1596        });
1597        let server = StubServer::start(state.clone());
1598        let client = SubsonicClient::new(&server.url(), "u", "p");
1599
1600        update_last_sync(&db, &server.url(), "u", 4_000_000_000).unwrap();
1601        let result = sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1602
1603        assert_eq!(
1604            result.albums_synced, 1,
1605            "an album with no usable timestamp must not be assumed old"
1606        );
1607    }
1608
1609    /// A server that lists songs for an empty query is synced in pages of
1610    /// songs, not one request per album.
1611    #[test]
1612    fn a_full_sync_pages_songs_instead_of_fetching_each_album() {
1613        let (db, _dir) = test_db();
1614        let state = Arc::new(StubState {
1615            albums: Mutex::new(stub_albums(1_234)),
1616            lists_songs: true,
1617            ..Default::default()
1618        });
1619        let server = StubServer::start(state.clone());
1620        let client = SubsonicClient::new(&server.url(), "u", "p");
1621
1622        let seen = Mutex::new(Vec::new());
1623        let result = sync_library(&db, &client, true, &server.url(), "u", &|p| {
1624            seen.lock().unwrap().push(p)
1625        })
1626        .unwrap();
1627
1628        assert!(state.album_calls.lock().unwrap().is_empty(), "no getAlbum");
1629        let mut pages = state.song_pages.lock().unwrap().clone();
1630        pages.sort();
1631        pages.dedup();
1632        assert_eq!(pages[..3], [0, 500, 1000]);
1633        assert_eq!(result.tracks_synced, 1_234);
1634        assert_eq!(result.albums_synced, 1_234);
1635        assert!(result.is_complete());
1636        assert!(get_last_sync(&db, &server.url()).unwrap().is_some());
1637
1638        let album: (Option<String>, Option<i32>) = db
1639            .conn
1640            .query_row(
1641                "SELECT al.remote_id, al.total_tracks FROM albums al WHERE al.title = 'Album 0007'",
1642                [],
1643                |r| Ok((r.get(0)?, r.get(1)?)),
1644            )
1645            .unwrap();
1646        assert_eq!(
1647            album,
1648            (Some("a0007".into()), Some(1)),
1649            "album metadata comes from the list"
1650        );
1651
1652        let seen = seen.into_inner().unwrap();
1653        let tracks: Vec<&SyncProgress> = seen
1654            .iter()
1655            .filter(|p| p.phase == SyncPhase::Tracks)
1656            .collect();
1657        assert!(tracks.len() >= 3, "progress at least per page");
1658        assert!(tracks.iter().all(|p| p.total == Some(1_234)));
1659        assert_eq!(tracks.last().unwrap().done, 1_234);
1660        assert_eq!(seen.last().unwrap().phase, SyncPhase::Finishing);
1661    }
1662
1663    /// Older servers answer an empty query with nothing. The album walk is
1664    /// what they get instead.
1665    #[test]
1666    fn a_server_that_lists_no_songs_is_synced_album_by_album() {
1667        let (db, _dir) = test_db();
1668        let state = Arc::new(StubState {
1669            albums: Mutex::new(stub_albums(3)),
1670            ..Default::default()
1671        });
1672        let server = StubServer::start(state.clone());
1673        let client = SubsonicClient::new(&server.url(), "u", "p");
1674
1675        let result = sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1676
1677        assert_eq!(state.song_pages.lock().unwrap().clone(), vec![0]);
1678        assert_eq!(state.album_calls.lock().unwrap().len(), 3);
1679        assert_eq!(result.tracks_synced, 3);
1680        assert!(result.is_complete());
1681    }
1682
1683    /// A page lost to the network leaves the run incomplete, so the watermark
1684    /// stays put and the next sync walks the library again.
1685    #[test]
1686    fn a_failed_song_page_does_not_advance_last_sync() {
1687        let (db, _dir) = test_db();
1688        let state = Arc::new(StubState {
1689            albums: Mutex::new(stub_albums(1_200)),
1690            lists_songs: true,
1691            failing_song_pages: Mutex::new([500].into_iter().collect()),
1692            ..Default::default()
1693        });
1694        let server = StubServer::start(state.clone());
1695        let client = SubsonicClient::new(&server.url(), "u", "p");
1696
1697        let result = sync_library(&db, &client, true, &server.url(), "u", &|_| {}).unwrap();
1698
1699        assert_eq!(result.pages_failed, 1);
1700        assert!(!result.is_complete());
1701        assert_eq!(get_last_sync(&db, &server.url()).unwrap(), None);
1702        let retries = state
1703            .song_pages
1704            .lock()
1705            .unwrap()
1706            .iter()
1707            .filter(|&&o| o == 500)
1708            .count();
1709        assert_eq!(retries as u32, PAGE_ATTEMPTS);
1710    }
1711
1712    /// An incremental sync fetches the few new albums one by one rather than
1713    /// walking every song.
1714    #[test]
1715    fn an_incremental_sync_does_not_walk_every_song() {
1716        let (db, _dir) = test_db();
1717        let state = Arc::new(StubState {
1718            albums: Mutex::new(stub_albums(3)),
1719            lists_songs: true,
1720            ..Default::default()
1721        });
1722        let server = StubServer::start(state.clone());
1723        let client = SubsonicClient::new(&server.url(), "u", "p");
1724        update_last_sync(&db, &server.url(), "u", 0).unwrap();
1725
1726        sync_library(&db, &client, false, &server.url(), "u", &|_| {}).unwrap();
1727
1728        assert!(state.song_pages.lock().unwrap().is_empty());
1729        assert_eq!(state.album_calls.lock().unwrap().len(), 3);
1730    }
1731
1732    #[test]
1733    fn sync_deduplicates_by_remote_id() {
1734        let (db, _dir) = test_db();
1735
1736        // First upsert.
1737        let meta1 = remote_track_meta("remote-dup", "Original Title", "Artist A", "Album X");
1738        let id1 = queries::upsert_track(&db.conn, &meta1).unwrap();
1739
1740        // Second upsert with same remote_id but different metadata.
1741        let meta2 = remote_track_meta("remote-dup", "Updated Title", "Artist A", "Album X");
1742        let id2 = queries::upsert_track(&db.conn, &meta2).unwrap();
1743
1744        // Should be the same row (dedup by remote_id).
1745        assert_eq!(id1, id2, "same remote_id should resolve to same track row");
1746
1747        // Verify the metadata was updated.
1748        let row = queries::get_track_row(&db.conn, id2)
1749            .unwrap()
1750            .expect("track should exist");
1751        assert_eq!(row.title, "Updated Title");
1752        assert_eq!(row.remote_id.as_deref(), Some("remote-dup"));
1753
1754        // Verify only one track exists.
1755        let stats = queries::library_stats(&db.conn).unwrap();
1756        assert_eq!(
1757            stats.total_tracks, 1,
1758            "should have exactly 1 track after dedup"
1759        );
1760    }
1761
1762    /// A koan server once published row ids and now publishes UUIDs; any
1763    /// server that rescans can renumber. A re-sync keeps the row, and with it
1764    /// the history and the favourite, rather than adding a second copy.
1765    #[test]
1766    fn resyncing_a_track_under_a_new_server_id_keeps_one_row() {
1767        let (db, _dir) = test_db();
1768        let before = remote_track_meta("46215", "Archangel", "Burial", "Untrue");
1769        let row = queries::upsert_synced_track(&db.conn, &before, &HashSet::from(["46215".into()]))
1770            .unwrap();
1771        queries::add_favourite(
1772            &db.conn,
1773            queries::LOCAL_USER,
1774            std::path::Path::new(before.remote_url.as_deref().unwrap()),
1775        )
1776        .unwrap();
1777
1778        let uid = "0199a0b2-7c4e-7d3a-9f1b-2c3d4e5f6a7b";
1779        let after = remote_track_meta(uid, "Archangel", "Burial", "Untrue");
1780        let again =
1781            queries::upsert_synced_track(&db.conn, &after, &HashSet::from([uid.into()])).unwrap();
1782
1783        assert_eq!(again, row, "the same row, not a second copy");
1784        assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 1);
1785        let (remote_id, album, artist): (String, String, String) = db
1786            .conn
1787            .query_row(
1788                "SELECT t.remote_id, al.remote_id, ar.remote_id FROM tracks t
1789                 JOIN albums al ON al.id = t.album_id JOIN artists ar ON ar.id = t.artist_id",
1790                [],
1791                |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
1792            )
1793            .unwrap();
1794        assert_eq!(remote_id, uid);
1795        assert_eq!(album, format!("album-of-{uid}"));
1796        assert_eq!(artist, format!("artist-of-{uid}"));
1797        let adopted: String = db
1798            .conn
1799            .query_row("SELECT uid FROM tracks WHERE id = ?1", [row], |r| r.get(0))
1800            .unwrap();
1801        assert_eq!(adopted, uid, "the row takes the server's uid");
1802        let favourites = queries::load_favourites(&db.conn, queries::LOCAL_USER).unwrap();
1803        assert_eq!(
1804            favourites,
1805            HashSet::from([std::path::PathBuf::from(after.remote_url.unwrap())]),
1806            "the favourite follows the new stream address"
1807        );
1808    }
1809
1810    /// Two entries a server lists with the same tags are two tracks, however
1811    /// alike: the second is not taken for the first under an old id.
1812    #[test]
1813    fn identical_entries_in_one_sync_stay_two_rows() {
1814        let (db, _dir) = test_db();
1815        let mut seen = HashSet::new();
1816        for id in ["dup-1", "dup-2"] {
1817            seen.insert(id.to_string());
1818            queries::upsert_synced_track(
1819                &db.conn,
1820                &remote_track_meta(id, "Archangel", "Burial", "Untrue"),
1821                &seen,
1822            )
1823            .unwrap();
1824        }
1825        assert_eq!(queries::library_stats(&db.conn).unwrap().total_tracks, 2);
1826    }
1827}