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