Skip to main content

koan_core/db/queries/
mod.rs

1mod albums;
2pub mod api_keys;
3mod artists;
4pub mod auth;
5pub mod batch;
6mod favourites;
7pub mod history;
8pub mod lyrics;
9pub mod playback_state;
10pub mod playlists;
11mod scan_cache;
12mod search;
13pub mod shares;
14pub(crate) mod sources;
15mod stats;
16pub mod tracks;
17pub mod uids;
18
19use std::path::PathBuf;
20
21// Re-exported so callers can `use queries::*`.
22pub use albums::*;
23pub use artists::*;
24pub use auth::LOCAL_USER;
25pub use batch::*;
26pub use favourites::*;
27pub use history::*;
28pub use lyrics::*;
29pub use playback_state::*;
30pub use playlists::*;
31pub use scan_cache::*;
32pub use search::*;
33pub use stats::*;
34pub use tracks::*;
35pub use uids::*;
36
37/// A write transaction that holds the write lock from its first statement.
38///
39/// What `Connection::unchecked_transaction` gives is deferred: it takes the
40/// lock at its first write, and if another connection wrote since its first
41/// read, SQLite refuses that upgrade outright rather than waiting, since
42/// waiting could deadlock. A scan chunk or a sync page that met another writer
43/// therefore failed every statement in it. Immediate waits for the lock the way
44/// any other statement does.
45pub fn write_transaction(
46    conn: &rusqlite::Connection,
47) -> rusqlite::Result<rusqlite::Transaction<'_>> {
48    rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)
49}
50
51/// Run `f` as one unit: its writes land together or not at all, readers never
52/// see them half done, and the write lock is taken once rather than per
53/// statement.
54///
55/// Inside a caller's transaction it is a savepoint, so it nests. Outside one it
56/// begins `IMMEDIATE`: a deferred transaction that reads before it writes
57/// fails outright, without waiting, when another connection wrote in between.
58pub fn atomically<T, E: From<rusqlite::Error>>(
59    conn: &rusqlite::Connection,
60    f: impl FnOnce() -> Result<T, E>,
61) -> Result<T, E> {
62    let (begin, commit, rollback) = if conn.is_autocommit() {
63        ("BEGIN IMMEDIATE", "COMMIT", "ROLLBACK")
64    } else {
65        (
66            "SAVEPOINT atomically",
67            "RELEASE atomically",
68            "ROLLBACK TO atomically; RELEASE atomically",
69        )
70    };
71    conn.execute_batch(begin)?;
72    let result = f();
73    match &result {
74        Ok(_) => conn.execute_batch(commit)?,
75        Err(_) => conn.execute_batch(rollback)?,
76    }
77    result
78}
79
80/// A list as one JSON array, for `IN (SELECT value FROM json_each(?))`.
81///
82/// One placeholder per item stops at SQLite's limit of 32,766 parameters, and
83/// a queue or a playlist may be longer than that. One parameter is not limited.
84pub fn json_list<T: serde::Serialize>(items: &[T]) -> String {
85    serde_json::to_string(items).unwrap_or_else(|_| "[]".into())
86}
87
88/// The half-open range of paths under a folder, for `path >= .0 AND path < .1`.
89///
90/// A prefix match on an indexed column, rather than `LIKE 'folder/%'` — which
91/// SQLite answers by reading every row, because a pattern is opaque to an
92/// index until it has been evaluated. It also takes the pattern out of the
93/// path: `LIKE` reads `_` as "any character" and folds ASCII case, so
94/// `/Volumes/My_Music` would match `/Volumes/My Music` and `/volumes/my_music`
95/// alike.
96///
97/// The trailing separator is what keeps `/Volumes/Music` out of
98/// `/Volumes/Music Backup`; the upper bound is the highest code point, so
99/// every path under the folder sorts below it.
100pub fn folder_prefix_range(folder: &std::path::Path) -> (String, String) {
101    let prefix = format!(
102        "{}{}",
103        folder
104            .to_string_lossy()
105            .trim_end_matches(std::path::MAIN_SEPARATOR),
106        std::path::MAIN_SEPARATOR
107    );
108    let upper = format!("{prefix}\u{10FFFF}");
109    (prefix, upper)
110}
111
112// --- Row types ---
113
114#[derive(Debug, Clone)]
115pub struct ArtistRow {
116    pub id: i64,
117    pub name: String,
118    pub sort_name: Option<String>,
119    pub remote_id: Option<String>,
120    /// Albums credited to this artist, and tracks across them. Aggregated in
121    /// the same query as the row itself — a count per artist would be one
122    /// query per row in a list thousands long.
123    pub album_count: i64,
124    pub track_count: i64,
125}
126
127#[derive(Debug, Clone)]
128pub struct AlbumRow {
129    pub id: i64,
130    pub title: String,
131    pub artist_id: i64,
132    pub artist_name: String,
133    pub date: Option<String>,
134    pub total_discs: Option<i32>,
135    pub total_tracks: Option<i32>,
136    pub codec: Option<String>,
137    pub label: Option<String>,
138    pub remote_id: Option<String>,
139    /// When the album entered the library — the server's `created` for remote
140    /// albums, otherwise the time it was first indexed.
141    pub added_at: Option<String>,
142}
143
144#[derive(Debug, Clone)]
145pub struct TrackRow {
146    pub id: i64,
147    pub album_id: Option<i64>,
148    pub artist_id: Option<i64>,
149    pub artist_name: String,
150    pub album_artist_name: String,
151    pub album_title: String,
152    pub disc: Option<i32>,
153    pub track_number: Option<i32>,
154    pub title: String,
155    pub duration_ms: Option<i64>,
156    pub path: Option<String>,
157    pub codec: Option<String>,
158    pub sample_rate: Option<i32>,
159    pub bit_depth: Option<i32>,
160    pub channels: Option<i32>,
161    pub bitrate: Option<i32>,
162    pub genre: Option<String>,
163    pub source: String,
164    pub remote_id: Option<String>,
165    pub cached_path: Option<String>,
166}
167
168/// Where to get audio data for playback. Local always wins.
169#[derive(Debug, Clone)]
170pub enum PlaybackSource {
171    Local(PathBuf),
172    Cached(PathBuf),
173    Remote(String),
174}
175
176#[derive(Debug, Clone, Default)]
177pub struct LibraryStats {
178    pub total_tracks: i64,
179    pub local_tracks: i64,
180    pub remote_tracks: i64,
181    pub cached_tracks: i64,
182    pub total_albums: i64,
183    pub total_artists: i64,
184}
185
186/// Metadata for inserting/updating a track.
187#[derive(Debug, Clone, PartialEq)]
188pub struct TrackMeta {
189    pub title: String,
190    pub artist: String,
191    pub album_artist: Option<String>,
192    pub album: String,
193    pub date: Option<String>,
194    pub disc: Option<i32>,
195    pub track_number: Option<i32>,
196    pub genre: Option<String>,
197    pub label: Option<String>,
198    pub duration_ms: Option<i64>,
199    pub codec: Option<String>,
200    pub sample_rate: Option<i32>,
201    pub bit_depth: Option<i32>,
202    pub channels: Option<i32>,
203    pub bitrate: Option<i32>,
204    pub size_bytes: Option<i64>,
205    pub mtime: Option<i64>,
206    pub path: Option<String>,
207    pub source: String,
208    pub remote_id: Option<String>,
209    pub remote_url: Option<String>,
210    /// The server's ids for the album and its artist.
211    ///
212    /// Carried alongside the track's own, because the server keys stars,
213    /// shares and cover art off them — a library synced without these has
214    /// albums and artists it can name but cannot refer to.
215    pub album_remote_id: Option<String>,
216    pub artist_remote_id: Option<String>,
217    /// MusicBrainz recording and release ids — `MUSICBRAINZ_TRACKID` and
218    /// `MUSICBRAINZ_ALBUMID` in a file's tags, `musicBrainzId` on a server's
219    /// song and album. Together they name one track whatever each source calls
220    /// the album; the recording alone recurs on every compilation it is on.
221    pub mbid: Option<String>,
222    pub album_mbid: Option<String>,
223    /// When the album this track belongs to entered the library. Remote sync
224    /// supplies the server's `created`; anything else leaves it and the album
225    /// is stamped with the time it was first seen.
226    pub album_added_at: Option<String>,
227}
228
229/// Test helper: build a sample TrackMeta for use in tests across sub-modules.
230#[cfg(test)]
231pub fn sample_meta(title: &str, artist: &str, album: &str) -> TrackMeta {
232    TrackMeta {
233        title: title.into(),
234        artist: artist.into(),
235        album_artist: Some(artist.into()),
236        album: album.into(),
237        date: Some("2024".into()),
238        disc: Some(1),
239        track_number: Some(1),
240        genre: Some("Electronic".into()),
241        label: None,
242        duration_ms: Some(240_000),
243        codec: Some("FLAC".into()),
244        sample_rate: Some(44100),
245        bit_depth: Some(16),
246        channels: Some(2),
247        bitrate: Some(1000),
248        size_bytes: Some(30_000_000),
249        mtime: Some(1700000000),
250        path: Some(format!("/music/{}/{}.flac", album, title)),
251        source: "local".into(),
252        remote_id: None,
253        album_remote_id: None,
254        artist_remote_id: None,
255        mbid: None,
256        album_mbid: None,
257        remote_url: None,
258        album_added_at: None,
259    }
260}
261
262#[cfg(test)]
263mod write_transaction_tests {
264    use std::time::Duration;
265
266    fn open(path: &std::path::Path) -> rusqlite::Connection {
267        let conn = rusqlite::Connection::open(path).unwrap();
268        conn.pragma_update(None, "journal_mode", "wal").unwrap();
269        conn.busy_timeout(Duration::from_secs(5)).unwrap();
270        conn
271    }
272
273    /// A deferred transaction that read before another connection committed
274    /// cannot write at all, however long the busy timeout; `write_transaction`
275    /// waits its turn instead. This is what a sync page meeting another writer
276    /// ran into.
277    #[test]
278    fn a_write_transaction_waits_where_a_deferred_one_fails() {
279        let dir = tempfile::tempdir().unwrap();
280        let path = dir.path().join("t.db");
281        open(&path)
282            .execute_batch("CREATE TABLE t (x INTEGER)")
283            .unwrap();
284
285        let deferred = open(&path);
286        let tx = deferred.unchecked_transaction().unwrap();
287        let _: i64 = tx
288            .query_row("SELECT COUNT(*) FROM t", [], |r| r.get(0))
289            .unwrap();
290        open(&path).execute("INSERT INTO t VALUES (1)", []).unwrap();
291        let err = tx.execute("INSERT INTO t VALUES (2)", []).unwrap_err();
292        assert!(
293            err.to_string().contains("locked") || err.to_string().contains("busy"),
294            "{err}"
295        );
296        drop(tx);
297
298        let path_for_writer = path.clone();
299        let writer = std::thread::spawn(move || {
300            let conn = open(&path_for_writer);
301            conn.execute_batch("BEGIN IMMEDIATE; INSERT INTO t VALUES (3)")
302                .unwrap();
303            std::thread::sleep(Duration::from_millis(200));
304            conn.execute_batch("COMMIT").unwrap();
305        });
306        std::thread::sleep(Duration::from_millis(50));
307        let conn = open(&path);
308        let tx = super::write_transaction(&conn).unwrap();
309        let _: i64 = tx
310            .query_row("SELECT COUNT(*) FROM t", [], |r| r.get(0))
311            .unwrap();
312        tx.execute("INSERT INTO t VALUES (4)", []).unwrap();
313        tx.commit().unwrap();
314        writer.join().unwrap();
315
316        let n: i64 = conn
317            .query_row("SELECT COUNT(*) FROM t", [], |r| r.get(0))
318            .unwrap();
319        assert_eq!(n, 3);
320    }
321}