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