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