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