Skip to main content

koan_core/remote/
sync.rs

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