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