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