Skip to main content

koan_server/graphql/
mod.rs

1mod helpers;
2mod jobs;
3mod loaders;
4mod mutations;
5mod queries;
6mod server;
7mod subscriptions;
8mod types;
9
10use std::ops::Deref;
11use std::path::PathBuf;
12use std::sync::Arc;
13use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
14use std::time::Duration;
15
16use async_graphql::dataloader::DataLoader;
17use async_graphql::{Context, Schema};
18use crossbeam_channel::{Sender, TrySendError};
19use koan_core::audio::viz::VizSnapshot;
20use koan_core::db::connection::Database;
21use koan_core::player::commands::PlayerCommand;
22use koan_core::player::state::{QueueItemId, SharedPlayerState};
23use uuid::Uuid;
24
25use koan_core::auth::Role;
26use loaders::DbLoader;
27use mutations::MutationRoot;
28use queries::QueryRoot;
29pub use server::{
30    ApiServerOpts, cmd_serve, cmd_serve_daemon, execute_in_process, start_api_background,
31};
32use subscriptions::SubscriptionRoot;
33
34use crate::auth::AuthUser;
35
36// ---------------------------------------------------------------------------
37// Connection pool
38// ---------------------------------------------------------------------------
39
40/// How long the dataloader gathers keys before running a batch.
41///
42/// The default 1ms closes the window while a wide selection set is still
43/// registering keys, splitting one query into several.
44const BATCH_WINDOW: Duration = Duration::from_millis(10);
45
46/// How long a resolver waits for a free connection before giving up. Long
47/// enough to ride out a scan chunk, short enough that a wedged writer surfaces
48/// as an error rather than a hung request.
49const ACQUIRE_TIMEOUT: Duration = Duration::from_secs(10);
50
51/// A fixed set of SQLite connections handed out per query, not per field.
52///
53/// `Database::open` runs the DDL batch, the migrations and a WAL checkpoint, so
54/// opening one per resolver field costs tens of statements before the actual
55/// query runs. Connections are opened lazily up to `max` and returned on drop.
56/// Deliberately not a single `Mutex<Connection>`: WAL gives concurrent readers,
57/// and one slow query must not block every other client.
58struct DbPool {
59    path: PathBuf,
60    /// Returned connections wait here. Bounded at `max`, so `try_send` on the
61    /// return path cannot block or fail for lack of room.
62    idle_tx: Sender<Database>,
63    idle_rx: crossbeam_channel::Receiver<Database>,
64    /// Connections in existence (idle plus checked out).
65    live: AtomicUsize,
66    /// Total connections ever opened. Instrumentation only — the fan-out tests
67    /// assert this stays bounded as a query's breadth grows.
68    opens: AtomicUsize,
69    /// Dataloader batches run. Instrumentation only — the N+1 tests assert one
70    /// batch serves a whole selection set.
71    batches: AtomicUsize,
72    /// Whether the schema and migrations have been applied to this path.
73    initialised: AtomicBool,
74    max: usize,
75}
76
77impl DbPool {
78    fn new(path: PathBuf) -> Arc<Self> {
79        // SQLite readers scale with cores; past that they only queue on the OS.
80        let max = std::thread::available_parallelism()
81            .map(|n| n.get())
82            .unwrap_or(4)
83            .clamp(4, 16);
84        let (idle_tx, idle_rx) = crossbeam_channel::bounded(max);
85        Arc::new(Self {
86            path,
87            idle_tx,
88            idle_rx,
89            live: AtomicUsize::new(0),
90            opens: AtomicUsize::new(0),
91            batches: AtomicUsize::new(0),
92            initialised: AtomicBool::new(false),
93            max,
94        })
95    }
96
97    /// The first connection applies the schema; every later one skips it.
98    fn open_new(&self) -> Result<Database, koan_core::db::connection::DbError> {
99        self.opens.fetch_add(1, Ordering::Relaxed);
100        if self.initialised.load(Ordering::Acquire) {
101            Database::open_existing(&self.path)
102        } else {
103            let db = Database::open(&self.path)?;
104            self.initialised.store(true, Ordering::Release);
105            Ok(db)
106        }
107    }
108
109    fn acquire(self: &Arc<Self>) -> async_graphql::Result<PooledDb> {
110        if let Ok(db) = self.idle_rx.try_recv() {
111            return Ok(PooledDb {
112                db: Some(db),
113                pool: self.clone(),
114            });
115        }
116
117        let mut live = self.live.load(Ordering::Relaxed);
118        while live < self.max {
119            match self.live.compare_exchange_weak(
120                live,
121                live + 1,
122                Ordering::AcqRel,
123                Ordering::Relaxed,
124            ) {
125                Ok(_) => {
126                    return match self.open_new() {
127                        Ok(db) => Ok(PooledDb {
128                            db: Some(db),
129                            pool: self.clone(),
130                        }),
131                        Err(e) => {
132                            self.live.fetch_sub(1, Ordering::AcqRel);
133                            Err(internal_error("db open", e))
134                        }
135                    };
136                }
137                Err(actual) => live = actual,
138            }
139        }
140
141        self.idle_rx
142            .recv_timeout(ACQUIRE_TIMEOUT)
143            .map(|db| PooledDb {
144                db: Some(db),
145                pool: self.clone(),
146            })
147            .map_err(|_| async_graphql::Error::new("database busy"))
148    }
149}
150
151/// A connection checked out of the pool, returned when the guard drops.
152struct PooledDb {
153    db: Option<Database>,
154    pool: Arc<DbPool>,
155}
156
157impl Deref for PooledDb {
158    type Target = Database;
159
160    fn deref(&self) -> &Database {
161        self.db.as_ref().expect("connection taken before drop")
162    }
163}
164
165impl Drop for PooledDb {
166    fn drop(&mut self) {
167        if let Some(db) = self.db.take()
168            && let Err(TrySendError::Full(_) | TrySendError::Disconnected(_)) =
169                self.pool.idle_tx.try_send(db)
170        {
171            self.pool.live.fetch_sub(1, Ordering::AcqRel);
172        }
173    }
174}
175
176// ---------------------------------------------------------------------------
177// DB handle wrapper (so we can put it in Context)
178// ---------------------------------------------------------------------------
179
180#[derive(Clone)]
181struct DbHandle {
182    pool: Arc<DbPool>,
183}
184
185impl DbHandle {
186    fn new(path: PathBuf) -> Self {
187        Self {
188            pool: DbPool::new(path),
189        }
190    }
191
192    fn acquire(&self) -> async_graphql::Result<PooledDb> {
193        self.pool.acquire()
194    }
195
196    /// A connection outside the pool, for work that runs for minutes and must
197    /// not deny a connection to request-path resolvers.
198    fn open_detached(&self) -> Result<Database, koan_core::db::connection::DbError> {
199        self.pool.open_new()
200    }
201
202    fn note_batch(&self) {
203        self.pool.batches.fetch_add(1, Ordering::Relaxed);
204    }
205
206    /// Connections opened since the schema was built.
207    #[cfg_attr(not(test), allow(dead_code))]
208    fn open_count(&self) -> usize {
209        self.pool.opens.load(Ordering::Relaxed)
210    }
211
212    /// Dataloader batches run since the schema was built.
213    #[cfg_attr(not(test), allow(dead_code))]
214    fn batch_count(&self) -> usize {
215        self.pool.batches.load(Ordering::Relaxed)
216    }
217}
218
219/// Run a rusqlite closure on the blocking pool with a pooled connection.
220///
221/// rusqlite is blocking start to finish. Calling it inline in an `async fn`
222/// parks a tokio worker for the duration, and enough of those at once starve
223/// the runtime — including the `ReaderStream`s feeding in-flight audio.
224async fn with_db<T, F>(ctx: &Context<'_>, f: F) -> async_graphql::Result<T>
225where
226    F: FnOnce(&Database) -> async_graphql::Result<T> + Send + 'static,
227    T: Send + 'static,
228{
229    let handle = ctx.data::<DbHandle>()?.clone();
230    blocking(move || {
231        let db = handle.acquire()?;
232        f(&db)
233    })
234    .await
235}
236
237/// Run any other blocking work (HTTP fetches, file decoding, tag reads) off the
238/// async workers.
239async fn blocking<T, F>(f: F) -> async_graphql::Result<T>
240where
241    F: FnOnce() -> async_graphql::Result<T> + Send + 'static,
242    T: Send + 'static,
243{
244    tokio::task::spawn_blocking(f)
245        .await
246        .map_err(|e| internal_error("blocking task", e))?
247}
248
249/// Log the detail, return a generic message.
250///
251/// SQLite errors quote the offending statement and filesystem errors quote
252/// absolute paths — a map of the host handed to whoever asked.
253pub(super) fn internal_error(context: &str, e: impl std::fmt::Display) -> async_graphql::Error {
254    log::error!("graphql {}: {}", context, e);
255    async_graphql::Error::new("internal error")
256}
257
258// ---------------------------------------------------------------------------
259// Schema builder
260// ---------------------------------------------------------------------------
261
262pub type KoanSchema = Schema<QueryRoot, MutationRoot, SubscriptionRoot>;
263
264pub fn build_schema(
265    state: Arc<SharedPlayerState>,
266    cmd_tx: Sender<PlayerCommand>,
267    db_path: PathBuf,
268    viz: Option<Arc<VizSnapshot>>,
269) -> KoanSchema {
270    build_schema_with(DbHandle::new(db_path), state, cmd_tx, viz)
271}
272
273fn build_schema_with(
274    handle: DbHandle,
275    state: Arc<SharedPlayerState>,
276    cmd_tx: Sender<PlayerCommand>,
277    viz: Option<Arc<VizSnapshot>>,
278) -> KoanSchema {
279    // Batching only, no caching: a schema-wide loader outlives the request, and
280    // favourites and library contents change under it.
281    let loader = DataLoader::new(DbLoader::new(handle.clone()), tokio::spawn).delay(BATCH_WINDOW);
282    loader.enable_all_cache(false);
283
284    let mut builder = Schema::build(QueryRoot, MutationRoot, SubscriptionRoot)
285        .data(handle)
286        .data(loader)
287        .data(jobs::JobRegistry::default())
288        .data(state)
289        .data(cmd_tx);
290    if let Some(viz) = viz {
291        builder = builder.data(viz);
292    }
293    // A single nested query can otherwise fan out across the whole library and
294    // pin the process for minutes.
295    builder.limit_depth(12).limit_complexity(2000).finish()
296}
297
298// ---------------------------------------------------------------------------
299// Shared helpers used by queries + mutations
300// ---------------------------------------------------------------------------
301
302fn parse_queue_item_id(s: &str) -> async_graphql::Result<QueueItemId> {
303    Uuid::parse_str(s)
304        .map(QueueItemId)
305        .map_err(|e| async_graphql::Error::new(format!("invalid queue item ID '{}': {}", s, e)))
306}
307
308/// How long to wait for room on the player command channel.
309///
310/// The channel is bounded(16) and the player can sit inside `start_playback`
311/// for about a second during a device sample-rate change. A blocking `send`
312/// there parks a tokio worker; this gives the player a moment to drain and
313/// then reports back rather than holding the thread.
314const CMD_SEND_TIMEOUT: Duration = Duration::from_millis(250);
315
316fn send_cmd(ctx: &Context<'_>, cmd: PlayerCommand) -> async_graphql::Result<()> {
317    let tx = ctx.data::<Sender<PlayerCommand>>()?;
318    send_cmd_via(tx, cmd)
319}
320
321fn send_cmd_via(tx: &Sender<PlayerCommand>, cmd: PlayerCommand) -> async_graphql::Result<()> {
322    tx.send_timeout(cmd, CMD_SEND_TIMEOUT)
323        .map_err(|_| async_graphql::Error::new("player busy — command not accepted"))
324}
325
326/// Extract the authenticated user from GraphQL context.
327/// Returns anonymous admin if no user is present (auth disabled or in-process).
328fn get_auth_user(ctx: &Context<'_>) -> AuthUser {
329    ctx.data::<AuthUser>()
330        .cloned()
331        .unwrap_or_else(|_| AuthUser::anonymous_admin())
332}
333
334/// Check that the current user has at least the required role.
335/// Returns an error suitable for GraphQL if the check fails.
336fn require_role(ctx: &Context<'_>, required: Role) -> async_graphql::Result<()> {
337    let user = get_auth_user(ctx);
338    if user.role.has_permission(required) {
339        Ok(())
340    } else {
341        Err(async_graphql::Error::new(format!(
342            "forbidden: requires {} role, you have {}",
343            required, user.role
344        )))
345    }
346}
347
348// ---------------------------------------------------------------------------
349// Tests
350// ---------------------------------------------------------------------------
351
352#[cfg(test)]
353mod tests {
354    use super::*;
355    use koan_core::db::connection::Database;
356    use koan_core::db::queries;
357    use koan_core::player::commands::CommandChannel;
358    use tempfile::TempDir;
359
360    /// One disposable configuration directory for the whole test binary.
361    ///
362    /// Deliberately not the per-test `TempDir`: enqueueing spawns downloads on
363    /// a thread that outlives the test which started it, and that thread reads
364    /// the configuration when it runs. Without this it reads the developer's
365    /// own — their library, their server, and their credentials.
366    fn isolate_config() {
367        static DIR: std::sync::OnceLock<TempDir> = std::sync::OnceLock::new();
368        koan_core::config::set_config_dir(DIR.get_or_init(|| TempDir::new().unwrap()).path());
369    }
370
371    fn test_schema() -> (
372        KoanSchema,
373        crossbeam_channel::Receiver<PlayerCommand>,
374        TempDir,
375    ) {
376        isolate_config();
377        let tmp = TempDir::new().unwrap();
378        let db_path = tmp.path().join("test.db");
379        let db = Database::open(&db_path).unwrap();
380        koan_core::db::schema::create_tables(&db.conn).unwrap();
381
382        let state = SharedPlayerState::new();
383        let ch = CommandChannel::new();
384        let tx = ch.tx.clone();
385        let rx = ch.rx.clone();
386
387        let schema = build_schema(state, tx, db_path, None);
388        (schema, rx, tmp)
389    }
390
391    /// Same schema, but keeping the `DbHandle` so tests can read its counters.
392    fn instrumented_schema() -> (
393        KoanSchema,
394        DbHandle,
395        crossbeam_channel::Receiver<PlayerCommand>,
396        TempDir,
397    ) {
398        let tmp = TempDir::new().unwrap();
399        let db_path = tmp.path().join("test.db");
400        let db = Database::open(&db_path).unwrap();
401        koan_core::db::schema::create_tables(&db.conn).unwrap();
402
403        let state = SharedPlayerState::new();
404        let ch = CommandChannel::new();
405        let handle = DbHandle::new(db_path);
406        let schema = build_schema_with(handle.clone(), state, ch.tx.clone(), None);
407        (schema, handle, ch.rx.clone(), tmp)
408    }
409
410    /// Seed a library on one connection — the per-track helper opens its own.
411    fn seed_library(db_path: &std::path::Path, artists: usize, albums: usize, tracks: usize) {
412        let db = Database::open(db_path).unwrap();
413        for a in 0..artists {
414            for al in 0..albums {
415                for t in 0..tracks {
416                    let meta = queries::TrackMeta {
417                        title: format!("Track {:03}-{}-{:03}", a, al, t),
418                        artist: format!("Artist {:03}", a),
419                        album_artist: Some(format!("Artist {:03}", a)),
420                        album: format!("Album {:03}-{}", a, al),
421                        track_number: Some(t as i32),
422                        disc: Some(1),
423                        date: Some("2024".into()),
424                        genre: Some("Electronic".into()),
425                        duration_ms: Some(240_000),
426                        path: Some(format!("/tmp/koan-test/{}/{}/{}.flac", a, al, t)),
427                        codec: Some("FLAC".into()),
428                        sample_rate: Some(44100),
429                        bit_depth: Some(16),
430                        channels: Some(2),
431                        bitrate: Some(1411),
432                        size_bytes: Some(42_000_000),
433                        mtime: Some(1700000000),
434                        source: "local".into(),
435                        remote_id: None,
436                        remote_url: None,
437                        album_remote_id: None,
438                        artist_remote_id: None,
439                        mbid: None,
440                        album_mbid: None,
441                        album_added_at: None,
442                        label: None,
443                    };
444                    queries::upsert_track(&db.conn, &meta).unwrap();
445                }
446            }
447        }
448    }
449
450    fn insert_test_track(db_path: &std::path::Path, title: &str, artist: &str, album: &str) -> i64 {
451        let db = Database::open(db_path).unwrap();
452        let meta = queries::TrackMeta {
453            title: title.to_string(),
454            artist: artist.to_string(),
455            album_artist: Some(artist.to_string()),
456            album: album.to_string(),
457            track_number: Some(1),
458            disc: Some(1),
459            date: Some("2024".into()),
460            genre: Some("Electronic".into()),
461            duration_ms: Some(240000),
462            path: Some(format!(
463                "/tmp/test/{}.flac",
464                title.to_lowercase().replace(' ', "_")
465            )),
466            codec: Some("FLAC".into()),
467            sample_rate: Some(44100),
468            bit_depth: Some(16),
469            channels: Some(2),
470            bitrate: Some(1411),
471            size_bytes: Some(42_000_000),
472            mtime: Some(1700000000),
473            source: "local".into(),
474            remote_id: None,
475            remote_url: None,
476            album_remote_id: None,
477            artist_remote_id: None,
478            mbid: None,
479            album_mbid: None,
480            album_added_at: None,
481            label: None,
482        };
483        queries::upsert_track(&db.conn, &meta).unwrap()
484    }
485
486    #[test]
487    fn schema_builds() {
488        let (_schema, _rx, _tmp) = test_schema();
489    }
490
491    #[tokio::test]
492    async fn library_stats_query() {
493        let (schema, _rx, tmp) = test_schema();
494        let db_path = tmp.path().join("test.db");
495        insert_test_track(&db_path, "Track1", "Artist1", "Album1");
496
497        let resp = schema
498            .execute("{ libraryStats { totalTracks totalAlbums totalArtists } }")
499            .await;
500        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
501        let data = resp.data.into_json().unwrap();
502        assert_eq!(data["libraryStats"]["totalTracks"], 1);
503        assert_eq!(data["libraryStats"]["totalAlbums"], 1);
504        assert_eq!(data["libraryStats"]["totalArtists"], 1);
505    }
506
507    #[tokio::test]
508    async fn artists_query() {
509        let (schema, _rx, tmp) = test_schema();
510        let db_path = tmp.path().join("test.db");
511        insert_test_track(&db_path, "T1", "Aphex Twin", "Drukqs");
512        insert_test_track(&db_path, "T2", "Boards of Canada", "MHTRTC");
513
514        let resp = schema
515            .execute("{ artists { edges { node { id name } } } }")
516            .await;
517        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
518        let data = resp.data.into_json().unwrap();
519        let edges = data["artists"]["edges"].as_array().unwrap();
520        assert_eq!(edges.len(), 2);
521    }
522
523    #[tokio::test]
524    async fn tracks_search() {
525        let (schema, _rx, tmp) = test_schema();
526        let db_path = tmp.path().join("test.db");
527        insert_test_track(&db_path, "Windowlicker", "Aphex Twin", "Windowlicker EP");
528        insert_test_track(&db_path, "Roygbiv", "Boards of Canada", "MHTRTC");
529
530        let resp = schema
531            .execute(r#"{ tracks(search: "Aphex") { edges { node { id title artist } } } }"#)
532            .await;
533        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
534        let data = resp.data.into_json().unwrap();
535        let edges = data["tracks"]["edges"].as_array().unwrap();
536        assert_eq!(edges.len(), 1);
537        assert_eq!(edges[0]["node"]["title"], "Windowlicker");
538    }
539
540    #[tokio::test]
541    async fn now_playing_stopped() {
542        let (schema, _rx, _tmp) = test_schema();
543        let resp = schema
544            .execute("{ nowPlaying { state positionMs track { title } } }")
545            .await;
546        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
547        let data = resp.data.into_json().unwrap();
548        assert_eq!(data["nowPlaying"]["state"], "STOPPED");
549    }
550
551    #[tokio::test]
552    async fn pause_mutation() {
553        let (schema, rx, _tmp) = test_schema();
554        let resp = schema.execute("mutation { pause { ok message } }").await;
555        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
556        let data = resp.data.into_json().unwrap();
557        assert_eq!(data["pause"]["ok"], true);
558        let cmd = rx.try_recv().unwrap();
559        assert!(matches!(cmd, PlayerCommand::Pause));
560    }
561
562    #[tokio::test]
563    async fn nested_artist_albums_tracks() {
564        let (schema, _rx, tmp) = test_schema();
565        let db_path = tmp.path().join("test.db");
566        insert_test_track(&db_path, "Vordhosbn", "Aphex Twin", "Drukqs");
567        insert_test_track(&db_path, "Avril 14th", "Aphex Twin", "Drukqs");
568
569        let resp = schema
570            .execute(
571                r#"{ artists(search: "Aphex") {
572                    edges { node {
573                        name
574                        albums { edges { node {
575                            title
576                            tracks { edges { node { title } } }
577                        } } }
578                    } }
579                } }"#,
580            )
581            .await;
582        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
583        let data = resp.data.into_json().unwrap();
584        let artist = &data["artists"]["edges"][0]["node"];
585        assert_eq!(artist["name"], "Aphex Twin");
586        let album = &artist["albums"]["edges"][0]["node"];
587        assert_eq!(album["title"], "Drukqs");
588        let tracks = album["tracks"]["edges"].as_array().unwrap();
589        assert_eq!(tracks.len(), 2);
590    }
591
592    #[tokio::test]
593    async fn pagination_has_next() {
594        let (schema, _rx, tmp) = test_schema();
595        let db_path = tmp.path().join("test.db");
596        for i in 0..5 {
597            insert_test_track(
598                &db_path,
599                &format!("Track{}", i),
600                "Artist",
601                &format!("Album{}", i),
602            );
603        }
604
605        let resp = schema
606            .execute(
607                r#"{ artists(first: 1) {
608                    edges { node { name } cursor }
609                    pageInfo { hasNextPage endCursor }
610                } }"#,
611            )
612            .await;
613        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
614        let data = resp.data.into_json().unwrap();
615        // Only 1 artist ("Artist"), so hasNextPage should be false
616        // since all 5 tracks are by the same artist.
617        assert_eq!(data["artists"]["edges"].as_array().unwrap().len(), 1);
618    }
619
620    #[tokio::test]
621    async fn clear_queue_mutation() {
622        let (schema, rx, _tmp) = test_schema();
623        let resp = schema
624            .execute("mutation { clearQueue { ok message } }")
625            .await;
626        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
627        let cmd = rx.try_recv().unwrap();
628        assert!(matches!(cmd, PlayerCommand::ClearPlaylist));
629    }
630
631    #[tokio::test]
632    async fn enqueue_mutation_adds_to_queue() {
633        let (schema, rx, tmp) = test_schema();
634        let db_path = tmp.path().join("test.db");
635
636        // Insert a track into the DB.
637        let track_id = insert_test_track(&db_path, "Windowlicker", "Aphex Twin", "Windowlicker EP");
638
639        // Execute the addToQueue mutation.
640        let query = format!(
641            "mutation {{ addToQueue(trackIds: [{}]) {{ ok message addedCount queueItemIds }} }}",
642            track_id
643        );
644        let resp = schema.execute(&query).await;
645        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
646
647        let data = resp.data.into_json().unwrap();
648        assert_eq!(data["addToQueue"]["ok"], true);
649        assert_eq!(data["addToQueue"]["addedCount"], 1);
650
651        let queue_ids = data["addToQueue"]["queueItemIds"].as_array().unwrap();
652        assert_eq!(queue_ids.len(), 1, "should return one queue item ID");
653
654        // Verify the PlayerCommand was sent through the channel.
655        // The mutation sends AddToPlaylist and then Play (auto-play when stopped).
656        let cmd = rx.try_recv().unwrap();
657        match cmd {
658            PlayerCommand::AddToPlaylist(items) => {
659                assert_eq!(items.len(), 1);
660                assert_eq!(items[0].title, "Windowlicker");
661                assert_eq!(items[0].artist, "Aphex Twin");
662                assert_eq!(items[0].album, "Windowlicker EP");
663            }
664            other => panic!("expected AddToPlaylist, got {:?}", other),
665        }
666
667        // Auto-play command should follow.
668        let play_cmd = rx.try_recv().unwrap();
669        assert!(
670            matches!(play_cmd, PlayerCommand::Play(_)),
671            "expected Play command for auto-play, got {:?}",
672            play_cmd
673        );
674    }
675
676    #[tokio::test]
677    async fn replace_queue_mutation_clears_and_enqueues() {
678        let (schema, rx, tmp) = test_schema();
679        let db_path = tmp.path().join("test.db");
680
681        let id1 = insert_test_track(&db_path, "Track A", "Artist", "Album");
682        let id2 = insert_test_track(&db_path, "Track B", "Artist", "Album");
683
684        let query = format!(
685            "mutation {{ replaceQueue(trackIds: [{}, {}]) {{ ok addedCount queueItemIds }} }}",
686            id1, id2
687        );
688        let resp = schema.execute(&query).await;
689        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
690
691        let data = resp.data.into_json().unwrap();
692        assert_eq!(data["replaceQueue"]["addedCount"], 2);
693
694        // One command, not clear-then-add-then-play: three commands down a
695        // bounded channel means the first track starts before the cursor lands
696        // on the one that was asked for.
697        match rx.try_recv().unwrap() {
698            PlayerCommand::ReplacePlaylist { items, start } => {
699                assert_eq!(items.len(), 2);
700                assert_eq!(start, 0, "defaults to the first track");
701            }
702            other => panic!("expected ReplacePlaylist, got {:?}", other),
703        }
704        assert!(rx.try_recv().is_err(), "no follow-up commands");
705    }
706
707    #[tokio::test]
708    async fn replace_queue_starts_where_it_was_asked_to() {
709        let (schema, rx, tmp) = test_schema();
710        let db_path = tmp.path().join("test.db");
711
712        let id1 = insert_test_track(&db_path, "Track A", "Artist", "Album");
713        let id2 = insert_test_track(&db_path, "Track B", "Artist", "Album");
714
715        let resp = schema
716            .execute(&format!(
717                "mutation {{ replaceQueue(trackIds: [{}, {}], startAt: 1) {{ ok }} }}",
718                id1, id2
719            ))
720            .await;
721        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
722
723        match rx.try_recv().unwrap() {
724            PlayerCommand::ReplacePlaylist { start, .. } => assert_eq!(start, 1),
725            other => panic!("expected ReplacePlaylist, got {:?}", other),
726        }
727    }
728
729    // ---- Playlists ----
730
731    /// The whole life of a playlist over the API: made, added to, reordered,
732    /// renamed, played, deleted.
733    #[tokio::test]
734    async fn playlists_round_trip_over_the_api() {
735        let (schema, rx, tmp) = test_schema();
736        let db_path = tmp.path().join("test.db");
737
738        let a = insert_test_track(&db_path, "Track A", "Artist", "Album");
739        let b = insert_test_track(&db_path, "Track B", "Artist", "Album");
740        let c = insert_test_track(&db_path, "Track C", "Artist", "Album");
741
742        let resp = schema
743            .execute(&format!(
744                "mutation {{ createPlaylist(name: \"Evening\", trackIds: [{a}, {b}]) \
745                 {{ id name trackCount }} }}"
746            ))
747            .await;
748        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
749        let data = resp.data.into_json().unwrap();
750        assert_eq!(data["createPlaylist"]["name"], "Evening");
751        assert_eq!(data["createPlaylist"]["trackCount"], 2);
752        let id = data["createPlaylist"]["id"].as_i64().unwrap();
753
754        let resp = schema
755            .execute(&format!(
756                "mutation {{ addToPlaylist(id: {id}, trackIds: [{c}]) {{ ok }} }}"
757            ))
758            .await;
759        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
760
761        let resp = schema
762            .execute(&format!("{{ playlistTracks(id: {id}) {{ title }} }}"))
763            .await;
764        let data = resp.data.into_json().unwrap();
765        let titles: Vec<&str> = data["playlistTracks"]
766            .as_array()
767            .unwrap()
768            .iter()
769            .map(|t| t["title"].as_str().unwrap())
770            .collect();
771        assert_eq!(
772            titles,
773            ["Track A", "Track B", "Track C"],
774            "playlist order is kept"
775        );
776
777        // A reorder is a wholesale replacement — the only shape Subsonic can
778        // express, so it is the only shape koan stores.
779        let resp = schema
780            .execute(&format!(
781                "mutation {{ setPlaylistTracks(id: {id}, trackIds: [{c}, {a}]) {{ ok }} }}"
782            ))
783            .await;
784        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
785
786        let resp = schema
787            .execute(&format!(
788                "mutation {{ renamePlaylist(id: {id}, name: \"Late\") {{ ok }} }}"
789            ))
790            .await;
791        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
792
793        let resp = schema.execute("{ playlists { id name trackCount } }").await;
794        let data = resp.data.into_json().unwrap();
795        assert_eq!(data["playlists"][0]["name"], "Late");
796        assert_eq!(data["playlists"][0]["trackCount"], 2);
797
798        // Playing it replaces the queue rather than adding to it.
799        let resp = schema
800            .execute(&format!("mutation {{ playPlaylist(id: {id}) {{ ok }} }}"))
801            .await;
802        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
803        match rx.try_recv().unwrap() {
804            PlayerCommand::ReplacePlaylist { items, start } => {
805                assert_eq!(items.len(), 2);
806                assert_eq!(start, 0);
807                assert!(items.iter().all(|i| i.playlist_entry_id.is_some()));
808            }
809            other => panic!("expected ReplacePlaylist, got {:?}", other),
810        }
811        assert!(rx.try_recv().is_err());
812
813        let resp = schema
814            .execute(&format!("mutation {{ deletePlaylist(id: {id}) {{ ok }} }}"))
815            .await;
816        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
817        let resp = schema.execute("{ playlists { id } }").await;
818        let data = resp.data.into_json().unwrap();
819        assert!(data["playlists"].as_array().unwrap().is_empty());
820    }
821
822    /// A queue item with no library row behind it cannot go in a playlist: a
823    /// playlist points at rows, not at paths.
824    #[tokio::test]
825    async fn saving_the_queue_keeps_only_what_the_library_knows() {
826        use koan_core::player::state::{ItemState, PlaylistItem};
827
828        let tmp = TempDir::new().unwrap();
829        let db_path = tmp.path().join("test.db");
830        let db = Database::open(&db_path).unwrap();
831        koan_core::db::schema::create_tables(&db.conn).unwrap();
832        drop(db);
833        let known = insert_test_track(&db_path, "Known", "Artist", "Album");
834
835        let state = SharedPlayerState::new();
836        let ch = CommandChannel::new();
837        let schema = build_schema(state.clone(), ch.tx.clone(), db_path, None);
838
839        let item = |title: &str, path: &str, db_id: Option<i64>| PlaylistItem {
840            playlist_entry_id: None,
841            id: QueueItemId::new(),
842            db_id,
843            path: std::path::PathBuf::from(path),
844            title: title.to_string(),
845            artist: "Artist".to_string(),
846            album_artist: "Artist".to_string(),
847            album: "Album".to_string(),
848            year: None,
849            codec: None,
850            track_number: None,
851            disc: None,
852            duration_ms: None,
853            state: ItemState::Ready,
854        };
855        state.add_items(vec![
856            item("Known", "/music/known.flac", Some(known)),
857            // Played from a file that was never indexed — a queue can hold one,
858            // a playlist cannot.
859            item("Stranger", "/elsewhere/stranger.flac", None),
860        ]);
861
862        let resp = schema
863            .execute("mutation { saveQueueAsPlaylist(name: \"Tonight\") { id trackCount } }")
864            .await;
865        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
866        let data = resp.data.into_json().unwrap();
867        assert_eq!(data["saveQueueAsPlaylist"]["trackCount"], 1);
868    }
869
870    // ---- Phase 1 tests: queue snapshot, viz, config, playlist version, subscriptions ----
871
872    #[tokio::test]
873    async fn queue_snapshot_has_version_and_status() {
874        let (schema, _rx, _tmp) = test_schema();
875
876        let resp = schema
877            .execute("{ queue { version entries { queueItemId status } hasPlaying queueCount } }")
878            .await;
879        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
880        let data = resp.data.into_json().unwrap();
881        // Empty queue.
882        assert_eq!(data["queue"]["version"], 0);
883        assert_eq!(data["queue"]["entries"].as_array().unwrap().len(), 0);
884        assert_eq!(data["queue"]["hasPlaying"], false);
885        assert_eq!(data["queue"]["queueCount"], 0);
886    }
887
888    #[tokio::test]
889    async fn queue_entries_have_status_and_download_progress() {
890        use koan_core::player::state::{ItemState, PlaylistItem};
891
892        // Build schema with a shared state we can manipulate directly.
893        let tmp = TempDir::new().unwrap();
894        let db_path = tmp.path().join("test.db");
895        let db = Database::open(&db_path).unwrap();
896        koan_core::db::schema::create_tables(&db.conn).unwrap();
897
898        let state = SharedPlayerState::new();
899        let ch = CommandChannel::new();
900        let schema = build_schema(state.clone(), ch.tx.clone(), db_path, None);
901
902        // Directly add items to the playlist (simulating what the player thread does).
903        let item = PlaylistItem {
904            playlist_entry_id: None,
905            id: QueueItemId::new(),
906            db_id: None,
907            path: std::path::PathBuf::from("/tmp/test/windowlicker.flac"),
908            title: "Windowlicker".to_string(),
909            artist: "Aphex Twin".to_string(),
910            album_artist: "Aphex Twin".to_string(),
911            album: "Windowlicker EP".to_string(),
912            year: None,
913            codec: Some("FLAC".to_string()),
914            track_number: Some(1),
915            disc: Some(1),
916            duration_ms: Some(240000),
917            state: ItemState::Ready,
918        };
919        state.add_items(vec![item]);
920
921        // Query the queue — should have one entry with QUEUED status.
922        let resp = schema
923            .execute(
924                "{ queue { version entries { queueItemId title status downloadProgress { downloaded total } isCurrent } finishedCount } }",
925            )
926            .await;
927        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
928        let data = resp.data.into_json().unwrap();
929        let entries = data["queue"]["entries"].as_array().unwrap();
930        assert_eq!(entries.len(), 1);
931        assert_eq!(entries[0]["title"], "Windowlicker");
932        // Without a cursor set, all entries are QUEUED.
933        assert_eq!(entries[0]["status"], "QUEUED");
934        assert_eq!(entries[0]["isCurrent"], false);
935        // Local track — no download progress.
936        assert!(entries[0]["downloadProgress"].is_null());
937    }
938
939    #[tokio::test]
940    async fn viz_frame_returns_none_without_viz() {
941        let (schema, _rx, _tmp) = test_schema();
942
943        let resp = schema
944            .execute("{ vizFrame { spectrum peaks vuLevels beatEnergy } }")
945            .await;
946        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
947        let data = resp.data.into_json().unwrap();
948        assert!(data["vizFrame"].is_null());
949    }
950
951    #[tokio::test]
952    async fn viz_frame_returns_data_with_viz() {
953        // Build schema with a VizSnapshot.
954        let tmp = TempDir::new().unwrap();
955        let db_path = tmp.path().join("test.db");
956        let db = Database::open(&db_path).unwrap();
957        koan_core::db::schema::create_tables(&db.conn).unwrap();
958
959        let state = SharedPlayerState::new();
960        let ch = CommandChannel::new();
961        let viz = koan_core::audio::viz::VizSnapshot::new();
962
963        // Write some test data.
964        let mut spectrum = [0.0f32; 48];
965        spectrum[0] = 0.75;
966        viz.write(koan_core::audio::viz::VizFrame {
967            spectrum,
968            peaks: [0.0; 48],
969            vu_levels: [0.42, 0.38],
970            beat_energy: 0.6,
971            timestamp: std::time::Instant::now(),
972            waveform: Vec::new(),
973        });
974
975        let schema = build_schema(state, ch.tx.clone(), db_path, Some(viz));
976
977        let resp = schema
978            .execute("{ vizFrame { spectrum peaks vuLevels beatEnergy waveform } }")
979            .await;
980        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
981        let data = resp.data.into_json().unwrap();
982        let frame = &data["vizFrame"];
983        assert!(!frame.is_null());
984        let spectrum = frame["spectrum"].as_array().unwrap();
985        assert_eq!(spectrum.len(), 48);
986        assert!((spectrum[0].as_f64().unwrap() - 0.75).abs() < 0.01);
987        let vu = frame["vuLevels"].as_array().unwrap();
988        assert_eq!(vu.len(), 2);
989        assert!((vu[0].as_f64().unwrap() - 0.42).abs() < 0.01);
990        assert!((frame["beatEnergy"].as_f64().unwrap() - 0.6).abs() < 0.01);
991        // Waveform empty — we didn't request includeWaveform.
992        let waveform = frame["waveform"].as_array().unwrap();
993        assert!(waveform.is_empty());
994    }
995
996    #[tokio::test]
997    async fn config_query() {
998        let (schema, _rx, _tmp) = test_schema();
999
1000        let resp = schema
1001            .execute(
1002                "{ config { libraryFolders replaygainMode targetFps artSize remoteEnabled graphqlPort } }",
1003            )
1004            .await;
1005        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1006        let data = resp.data.into_json().unwrap();
1007        let cfg = &data["config"];
1008        // Defaults from Config::default().
1009        assert!(cfg["libraryFolders"].is_array());
1010        assert!(cfg["targetFps"].as_i64().unwrap() > 0);
1011        assert!(cfg["artSize"].as_i64().unwrap() > 0);
1012    }
1013
1014    #[tokio::test]
1015    async fn playlist_version_query() {
1016        let (schema, _rx, _tmp) = test_schema();
1017
1018        let resp = schema.execute("{ playlistVersion }").await;
1019        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1020        let data = resp.data.into_json().unwrap();
1021        assert_eq!(data["playlistVersion"], 0);
1022    }
1023
1024    #[tokio::test]
1025    async fn subscription_types_in_schema() {
1026        // Verify that subscriptions are registered by introspecting the schema.
1027        let (schema, _rx, _tmp) = test_schema();
1028
1029        let resp = schema
1030            .execute("{ __schema { subscriptionType { fields { name } } } }")
1031            .await;
1032        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1033        let data = resp.data.into_json().unwrap();
1034        let fields = data["__schema"]["subscriptionType"]["fields"]
1035            .as_array()
1036            .unwrap();
1037        let names: Vec<&str> = fields.iter().filter_map(|f| f["name"].as_str()).collect();
1038        assert!(
1039            names.contains(&"nowPlaying"),
1040            "missing nowPlaying subscription"
1041        );
1042        assert!(
1043            names.contains(&"queueUpdated"),
1044            "missing queueUpdated subscription"
1045        );
1046        assert!(names.contains(&"vizFrame"), "missing vizFrame subscription");
1047    }
1048    // -- Fan-out and pagination --
1049
1050    #[tokio::test]
1051    async fn nested_fan_out_opens_a_bounded_number_of_connections() {
1052        let (schema, handle, _rx, tmp) = instrumented_schema();
1053        seed_library(&tmp.path().join("test.db"), 10, 3, 4);
1054
1055        let before = handle.open_count();
1056        let resp = schema
1057            .execute(
1058                "{ artists(first: 10) { edges { node { name albumCount trackCount \
1059                   albums(first: 10) { edges { node { title trackCount totalDurationMs \
1060                   tracks(first: 10) { edges { node { title isFavourite } } } } } } } } } }",
1061            )
1062            .await;
1063        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1064
1065        let data = resp.data.into_json().unwrap();
1066        let artists = data["artists"]["edges"].as_array().unwrap();
1067        assert_eq!(artists.len(), 10, "query did not actually fan out");
1068        assert_eq!(artists[0]["node"]["albumCount"], 3);
1069        assert_eq!(artists[0]["node"]["trackCount"], 12);
1070
1071        // Before pooling this was one full `Database::open` — DDL batch,
1072        // migrations and a WAL checkpoint — per resolver field.
1073        let opened = handle.open_count() - before;
1074        assert!(opened <= 16, "opened {} connections", opened);
1075    }
1076
1077    #[tokio::test]
1078    async fn is_favourite_batches_into_one_query() {
1079        let (schema, handle, _rx, tmp) = instrumented_schema();
1080        seed_library(&tmp.path().join("test.db"), 1, 1, 100);
1081
1082        let before = handle.batch_count();
1083        let resp = schema
1084            .execute("{ tracks(first: 100) { edges { node { title isFavourite } } } }")
1085            .await;
1086        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1087        let data = resp.data.into_json().unwrap();
1088        assert_eq!(data["tracks"]["edges"].as_array().unwrap().len(), 100);
1089
1090        // The invariant is that the query count does not scale with the row
1091        // count: without the dataloader this was one full `favourites` scan per
1092        // track. The dataloader's gather window can close more than once under
1093        // load, so assert the property rather than an exact batch count.
1094        let batches = handle.batch_count() - before;
1095        assert!(
1096            batches < 10,
1097            "{batches} batches for 100 tracks — expected a handful, not one per row"
1098        );
1099    }
1100
1101    #[tokio::test]
1102    async fn tracks_without_first_returns_the_default_page() {
1103        let (schema, _handle, _rx, tmp) = instrumented_schema();
1104        seed_library(
1105            &tmp.path().join("test.db"),
1106            1,
1107            1,
1108            helpers::DEFAULT_PAGE + 20,
1109        );
1110
1111        let resp = schema
1112            .execute("{ tracks { edges { node { title } } pageInfo { hasNextPage } } }")
1113            .await;
1114        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1115        let data = resp.data.into_json().unwrap();
1116        assert_eq!(
1117            data["tracks"]["edges"].as_array().unwrap().len(),
1118            helpers::DEFAULT_PAGE
1119        );
1120        assert_eq!(data["tracks"]["pageInfo"]["hasNextPage"], true);
1121    }
1122
1123    #[tokio::test]
1124    async fn first_is_clamped_to_the_maximum_page() {
1125        let (schema, _handle, _rx, tmp) = instrumented_schema();
1126        seed_library(&tmp.path().join("test.db"), 1, 1, 10);
1127
1128        let resp = schema
1129            .execute("{ tracks(first: 100000) { edges { node { title } } } }")
1130            .await;
1131        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1132        let data = resp.data.into_json().unwrap();
1133        assert_eq!(data["tracks"]["edges"].as_array().unwrap().len(), 10);
1134    }
1135
1136    #[tokio::test]
1137    async fn negative_first_is_empty_not_a_panic() {
1138        let (schema, _handle, _rx, tmp) = instrumented_schema();
1139        seed_library(&tmp.path().join("test.db"), 2, 1, 3);
1140
1141        for query in [
1142            "{ tracks(first: -1) { edges { node { title } } } }",
1143            "{ artists(first: -1) { edges { node { name } } } }",
1144            "{ albums(first: -1) { edges { node { title } } } }",
1145        ] {
1146            let resp = schema.execute(query).await;
1147            assert!(resp.errors.is_empty(), "{}: {:?}", query, resp.errors);
1148        }
1149    }
1150
1151    #[tokio::test]
1152    async fn sort_arguments_are_honoured() {
1153        let (schema, _handle, _rx, tmp) = instrumented_schema();
1154        seed_library(&tmp.path().join("test.db"), 3, 1, 2);
1155
1156        let resp = schema
1157            .execute(
1158                "{ tracks(first: 100, sortBy: TITLE, sortDir: DESC) { edges { node { title } } } }",
1159            )
1160            .await;
1161        assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1162        let data = resp.data.into_json().unwrap();
1163        let titles: Vec<String> = data["tracks"]["edges"]
1164            .as_array()
1165            .unwrap()
1166            .iter()
1167            .map(|e| e["node"]["title"].as_str().unwrap().to_string())
1168            .collect();
1169        let mut sorted = titles.clone();
1170        sorted.sort();
1171        sorted.reverse();
1172        assert_eq!(titles, sorted, "sortBy/sortDir were ignored");
1173    }
1174
1175    /// A blocking resolver must not hold the runtime. On a single-threaded
1176    /// runtime, anything running inline would freeze every other task for the
1177    /// whole query — which is what stalled in-flight audio streams.
1178    #[test]
1179    fn a_slow_query_does_not_stall_other_tasks() {
1180        use std::sync::atomic::{AtomicUsize, Ordering};
1181
1182        let rt = tokio::runtime::Builder::new_current_thread()
1183            .enable_time()
1184            .build()
1185            .unwrap();
1186        let (schema, _handle, _rx, tmp) = instrumented_schema();
1187        seed_library(&tmp.path().join("test.db"), 20, 5, 8);
1188
1189        rt.block_on(async move {
1190            let ticks = Arc::new(AtomicUsize::new(0));
1191            let counter = ticks.clone();
1192            let ticker = tokio::spawn(async move {
1193                loop {
1194                    tokio::time::sleep(Duration::from_millis(1)).await;
1195                    counter.fetch_add(1, Ordering::Relaxed);
1196                }
1197            });
1198
1199            // fuzzySearch loads the library and builds a whole Nucleo matcher.
1200            let mut running = Vec::new();
1201            for _ in 0..4 {
1202                let schema = schema.clone();
1203                running.push(tokio::spawn(async move {
1204                    schema
1205                        .execute("{ fuzzySearch(query: \"track\") { id } }")
1206                        .await
1207                }));
1208            }
1209            for handle in running {
1210                let resp = handle.await.unwrap();
1211                assert!(resp.errors.is_empty(), "errors: {:?}", resp.errors);
1212            }
1213
1214            ticker.abort();
1215            assert!(
1216                ticks.load(Ordering::Relaxed) >= 2,
1217                "no other task ran while the queries were in flight"
1218            );
1219        });
1220    }
1221}