koan-server 0.37.0

GraphQL, Subsonic REST, and MCP server for koan music player.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
//! MCP (Model Context Protocol) server for koan.
//!
//! Exposes the GraphQL schema as MCP tools for Claude Desktop / MCP clients.

use std::path::PathBuf;
use std::sync::Arc;

use crossbeam_channel::Sender;
use koan_core::player::commands::PlayerCommand;
use koan_core::player::state::SharedPlayerState;
use rmcp::handler::server::router::tool::ToolRouter;
use rmcp::handler::server::wrapper::Json;
use rmcp::model::{ServerCapabilities, ServerConfig};
use rmcp::{ServerHandler, schemars, tool_router};
use serde::{Deserialize, Serialize};

// ---------------------------------------------------------------------------
// Parameter types
// ---------------------------------------------------------------------------

#[derive(Debug, Deserialize, schemars::JsonSchema)]
pub struct GraphqlParams {
    #[schemars(
        description = "GraphQL query or mutation string. Use the schema_sdl tool first to learn available types, queries, mutations, and filter parameters."
    )]
    pub query: String,
    #[schemars(description = "Optional JSON object of query variables")]
    pub variables: Option<serde_json::Value>,
}

// ---------------------------------------------------------------------------
// Response types
// ---------------------------------------------------------------------------

/// GraphQL execution result wrapper — MCP spec requires outputSchema to be an object type.
#[derive(Debug, Serialize, schemars::JsonSchema)]
pub struct GraphqlResponse {
    /// The GraphQL response JSON (contains data and/or errors fields).
    pub result: serde_json::Value,
}

// ---------------------------------------------------------------------------
// MCP Server
// ---------------------------------------------------------------------------

/// A koan account, sent by the MCP gateway with every request once the
/// gateway has signed its user in. Only the HTTP transport reads them, and it
/// is reachable from the gateway alone.
pub const USERNAME_HEADER: &str = "x-koan-username";
pub const PASSWORD_HEADER: &str = "x-koan-password";

#[derive(Clone)]
pub struct KoanMcpServer {
    #[allow(dead_code)]
    tool_router: ToolRouter<Self>,
    graphql_schema: crate::graphql::KoanSchema,
    /// Checks account headers; `None` on stdio, which has no headers.
    users: Option<Arc<crate::auth::password::PasswordVerifier>>,
}

impl KoanMcpServer {
    pub fn new(
        state: Arc<SharedPlayerState>,
        cmd_tx: Sender<PlayerCommand>,
        db_path: PathBuf,
    ) -> Self {
        let graphql_schema =
            crate::graphql::build_schema(state.clone(), cmd_tx.clone(), db_path.clone(), None);
        Self {
            tool_router: Self::tool_router(),
            graphql_schema,
            users: None,
        }
    }

    /// The role a request acts with: the account in its headers, or
    /// `mcp_role()` when it names none. With `KOAN_MCP_REQUIRE_LOGIN=1`, a
    /// request naming no account is refused.
    fn role(&self, extensions: &rmcp::model::Extensions) -> Result<koan_core::auth::Role, String> {
        let Some(users) = &self.users else {
            return Ok(mcp_role());
        };
        let parts = extensions.get::<axum::http::request::Parts>();
        let get = |h: &str| {
            parts
                .and_then(|p| p.headers.get(h))
                .and_then(|v| v.to_str().ok())
                .map(str::trim)
                .filter(|v| !v.is_empty())
        };
        match (get(USERNAME_HEADER), get(PASSWORD_HEADER)) {
            (Some(u), Some(p)) => users
                .verify(u, p)
                .ok_or_else(|| "kōan rejected that username and password".to_string()),
            _ if std::env::var("KOAN_MCP_REQUIRE_LOGIN").is_ok_and(|v| v == "1") => Err(format!(
                "this kōan needs an account: send {USERNAME_HEADER} and {PASSWORD_HEADER}"
            )),
            _ => Ok(mcp_role()),
        }
    }
}

use rmcp::handler::server::wrapper::Parameters;
use rmcp::tool;

/// Role the MCP `graphql` tool executes at.
///
/// The transport carries no credential, so anything reachable here is reachable
/// by whoever can talk to the MCP process. `User` covers everything the tool
/// advertises — browsing, playback, queue, favourites, playlists, radio — and
/// leaves out the admin mutations that move files on disk (`organize*`), rewrite
/// config, or change the output device. `KOAN_MCP_ADMIN=1` opts back in.
fn mcp_role() -> koan_core::auth::Role {
    if std::env::var("KOAN_MCP_ADMIN").is_ok_and(|v| v == "1") {
        koan_core::auth::Role::Admin
    } else {
        koan_core::auth::Role::User
    }
}

#[tool_router]
impl KoanMcpServer {
    #[tool(
        description = "Get the full GraphQL schema in SDL format. CALL THIS FIRST to learn all \
        available queries, mutations, types, and filter parameters. The schema is the complete \
        reference for everything koan can do — library discovery, playback control, queue \
        management, favourites, playlists, radio mode, device switching, and more."
    )]
    fn schema_sdl(&self) -> Json<GraphqlResponse> {
        let sdl = self.graphql_schema.sdl();
        Json(GraphqlResponse {
            result: serde_json::Value::String(sdl),
        })
    }

    #[tool(
        description = "Execute a GraphQL query or mutation against the koan music player. \
        This is the primary interface for ALL operations — library browsing, playback control, \
        queue management, favourites, playlists, radio, devices.\n\n\
        Call schema_sdl first to learn the full schema.\n\n\
        Quick examples:\n\
        - Search: { tracks(search: \"aphex\") { edges { node { id title artist album } } } }\n\
        - Filter: { albums(yearEnd: 1995, codec: \"FLAC\") { edges { node { title artistName date } } } }\n\
        - Now playing: { nowPlaying { state positionMs track { title artist codec sampleRate } } }\n\
        - Queue tracks: mutation { addToQueue(trackIds: [42, 43]) { ok addedCount } }\n\
        - Play/pause: mutation { pause { ok } } / mutation { resume { ok } }\n\
        - Playlist: mutation { saveQueueAsPlaylist(name: \"techno\") { id name } }\n\
        - Radio: mutation { enableRadio { ok } }\n\n\
        Track IDs are integers from the library. Queue item IDs are UUIDs from the queue.\n\
        All string filters are case-insensitive substrings."
    )]
    fn graphql(
        &self,
        Parameters(params): Parameters<GraphqlParams>,
        extensions: rmcp::model::Extensions,
    ) -> Result<Json<GraphqlResponse>, String> {
        let schema = self.graphql_schema.clone();
        let query = params.query;
        let variables = params.variables;
        let rt =
            tokio::runtime::Handle::try_current().map_err(|_| "no tokio runtime".to_string())?;
        // Inside block_in_place too: a first sign-in runs argon2.
        let result = tokio::task::block_in_place(|| {
            let role = self.role(&extensions)?;
            Ok::<_, String>(rt.block_on(crate::graphql::execute_in_process(
                &schema, &query, variables, role,
            )))
        })?;
        Ok(Json(GraphqlResponse { result }))
    }
}

#[rmcp::tool_handler]
impl ServerHandler for KoanMcpServer {
    fn get_info(&self) -> ServerConfig {
        ServerConfig::new(ServerCapabilities::builder().enable_tools().build())
            .with_server_info(rmcp::model::Implementation::new(
                "koan",
                env!("CARGO_PKG_VERSION"),
            ))
            .with_instructions(
                "koan is a bit-perfect music player. You control it entirely via GraphQL.\n\n\
             ## How to use\n\
             1. Call `schema_sdl` to get the full GraphQL schema\n\
             2. Use the `graphql` tool for ALL queries and mutations\n\n\
             ## What you can do\n\
             - **Discover music**: query `artists`, `albums`, `tracks` with rich filters \
               (genre, year range, codec, sample rate, bit depth, duration, favourites)\n\
             - **Control playback**: mutations `play`, `pause`, `resume`, `stop`, `next`, \
               `previous`, `seek`\n\
             - **Manage queue**: `addToQueue`, `replaceQueue`, `removeFromQueue`, `moveInQueue`, \
               `clearQueue`, `undo`, `redo`\n\
             - **Favourites**: `favourite`, `unfavourite`, `toggleFavourite` (auto-syncs to \
               Subsonic/Navidrome). Filter any query with `favouritesOnly: true`\n\
             - **Playlists**: query `playlists`/`playlistTracks`; `createPlaylist`, \
               `saveQueueAsPlaylist`, `addToPlaylist`, `setPlaylistTracks`, `renamePlaylist`, \
               `deletePlaylist`, `playPlaylist`. Synced to Subsonic/Navidrome\n\
             - **Radio**: `enableRadio`, `disableRadio` — auto-queues similar tracks\n\
             - **Play on the user's phone or Mac**: query `clients` for the koan apps \
               linked to this server, then `playOnClient(trackIds, client)` to replace its \
               queue (or `enqueue: true` to append), and `controlClient` to pause, resume or \
               skip. Build the list with the library queries first; the music plays on that \
               device, not on the server\n\
             - **Devices**: query `devices`; `setDevice`/`clearDevice` need `KOAN_MCP_ADMIN=1`\n\
             - **History**: query `playHistory`, `similarArtists`\n\
             - **Sharing**: `createShare(trackIds, description)` returns a public link anyone can \
               open without an account; query `shares` to list them, `updateShare` to set an \
               expiry, `deleteShare` to revoke one. A link is public, so confirm with the user \
               before making one\n\n\
             ## Not available\n\
             Admin mutations — `organize*` (moves files on disk), `updateConfig`, \
             `triggerScan` — are refused unless `KOAN_MCP_ADMIN=1` is set.\n\n\
             ## ID conventions\n\
             - Track IDs: integers from the library database\n\
             - Queue item IDs: UUIDs assigned when tracks enter the queue",
            )
    }
}

/// Serve MCP over streamable HTTP at `addr`/mcp, on a thread of its own.
///
/// The gateway authenticates its user and then forwards a koan account in
/// `x-koan-username` / `x-koan-password`, which set the role each request acts
/// with. The headers are trusted to come from the gateway, so this listener
/// must not be reachable from anywhere else: bind it to an address only the
/// gateway can reach, and keep it off the public GraphQL port.
pub fn spawn_http(
    addr: std::net::SocketAddr,
    state: Arc<SharedPlayerState>,
    cmd_tx: Sender<PlayerCommand>,
    pool: Arc<koan_core::db::pool::Pool>,
) -> std::io::Result<std::thread::JoinHandle<()>> {
    use rmcp::transport::streamable_http_server::{
        StreamableHttpServerConfig, StreamableHttpService, session::local::LocalSessionManager,
    };
    let mut template = KoanMcpServer::new(state, cmd_tx, pool.path().to_path_buf());
    template.users = Some(Arc::new(crate::auth::password::PasswordVerifier::new(pool)));
    // Bound here rather than on the thread, so a taken port fails the start.
    let listener = std::net::TcpListener::bind(addr)?;
    listener.set_nonblocking(true)?;
    std::thread::Builder::new()
        .name("koan-mcp-http".into())
        .spawn(move || {
            let rt = tokio::runtime::Builder::new_multi_thread()
                .enable_all()
                .build()
                .expect("failed to create tokio runtime");
            rt.block_on(async move {
                let service = StreamableHttpService::new(
                    move || Ok(template.clone()),
                    Arc::new(LocalSessionManager::default()),
                    // The gateway forwards its own Host header, which rmcp's
                    // DNS-rebinding allowlist would refuse; reachability is
                    // what guards this listener.
                    StreamableHttpServerConfig::default().disable_allowed_hosts(),
                );
                let app = axum::Router::new().nest_service("/mcp", service);
                let listener =
                    tokio::net::TcpListener::from_std(listener).expect("listener from std");
                if let Err(e) = axum::serve(listener, app).await {
                    log::error!("MCP HTTP server stopped: {e}");
                }
            });
        })
}

/// Entry point for `koan mcp` — starts a headless player with an MCP server on stdio.
pub fn cmd_mcp() {
    use koan_core::player::Player;
    use rmcp::ServiceExt;

    // Validate DB is accessible before starting the server.
    let _db = koan_core::db::connection::Database::open_default().expect("failed to open database");
    let db_path = koan_core::config::db_path();

    // Spawn the player engine (headless — no TUI).
    let (state, _timeline, _viz, cmd_tx) = Player::spawn();

    let server = KoanMcpServer::new(state, cmd_tx, db_path);

    // Run the MCP server on the tokio runtime (blocking the main thread).
    let rt = tokio::runtime::Runtime::new().expect("failed to create tokio runtime");
    rt.block_on(async {
        let transport = rmcp::transport::io::stdio();
        let service = server
            .serve(transport)
            .await
            .expect("failed to start MCP server");
        let _ = service.waiting().await;
    });
}

// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------

#[cfg(test)]
mod tests {
    use super::*;
    use koan_core::db::connection::Database;
    use koan_core::db::queries;
    use koan_core::player::commands::CommandChannel;
    use tempfile::TempDir;

    fn test_server() -> (KoanMcpServer, CommandChannel, TempDir) {
        let tmp = TempDir::new().unwrap();
        let db_path = tmp.path().join("test.db");
        let db = Database::open(&db_path).unwrap();
        koan_core::db::schema::create_tables(&db.conn).unwrap();

        let state = SharedPlayerState::new();
        let ch = CommandChannel::new();
        let tx = ch.tx.clone();

        let server = KoanMcpServer::new(state, tx, db_path);
        (server, ch, tmp)
    }

    fn with_headers(headers: &[(&str, &str)]) -> rmcp::model::Extensions {
        let mut req = axum::http::Request::builder();
        for (k, v) in headers {
            req = req.header(*k, *v);
        }
        let (parts, ()) = req.body(()).unwrap().into_parts();
        let mut ext = rmcp::model::Extensions::new();
        ext.insert(parts);
        ext
    }

    #[test]
    fn gateway_headers_act_as_that_account() {
        use koan_core::auth::Role;
        let (mut server, _ch, tmp) = test_server();
        let db_path = tmp.path().join("test.db");
        let db = Database::open(&db_path).unwrap();
        queries::auth::create_user(&db.conn, "owner", "sesame", Role::Admin).unwrap();
        queries::auth::create_user(&db.conn, "mate", "hunter22", Role::Readonly).unwrap();
        server.users = Some(Arc::new(crate::auth::password::PasswordVerifier::new(
            Arc::new(koan_core::db::pool::Pool::new(db_path)),
        )));

        let as_ = |u: &str, p: &str| {
            server.role(&with_headers(&[(USERNAME_HEADER, u), (PASSWORD_HEADER, p)]))
        };
        assert_eq!(as_("owner", "sesame"), Ok(Role::Admin));
        assert_eq!(as_("mate", "hunter22"), Ok(Role::Readonly));
        assert!(as_("owner", "wrong").is_err());
        // No account named: the transport's default role.
        assert_eq!(server.role(&with_headers(&[])), Ok(mcp_role()));
    }

    fn insert_test_track(db_path: &std::path::Path, title: &str, artist: &str, album: &str) -> i64 {
        let db = Database::open(db_path).unwrap();
        let meta = queries::TrackMeta {
            title: title.to_string(),
            artist: artist.to_string(),
            album_artist: Some(artist.to_string()),
            album: album.to_string(),
            track_number: Some(1),
            disc: Some(1),
            date: Some("2024".into()),
            genre: Some("Electronic".into()),
            duration_ms: Some(240000),
            path: Some(format!(
                "/tmp/test/{}.flac",
                title.to_lowercase().replace(' ', "_")
            )),
            codec: Some("FLAC".into()),
            sample_rate: Some(44100),
            bit_depth: Some(16),
            channels: Some(2),
            bitrate: Some(1411),
            size_bytes: Some(42_000_000),
            mtime: Some(1700000000),
            source: "local".into(),
            remote_id: None,
            remote_url: None,
            album_remote_id: None,
            artist_remote_id: None,
            mbid: None,
            album_mbid: None,
            album_added_at: None,
            label: None,
        };
        queries::upsert_track(&db.conn, &meta).unwrap()
    }

    #[test]
    fn schema_sdl_returns_schema() {
        let (server, _ch, _tmp) = test_server();
        let Json(resp) = server.schema_sdl();
        let sdl = resp.result.as_str().unwrap();
        assert!(sdl.contains("type QueryRoot"));
        assert!(sdl.contains("type MutationRoot"));
        assert!(sdl.contains("artists"));
        assert!(sdl.contains("nowPlaying"));
    }

    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
    async fn graphql_query_works() {
        let (server, _ch, tmp) = test_server();
        let db_path = tmp.path().join("test.db");
        insert_test_track(&db_path, "Windowlicker", "Aphex Twin", "Windowlicker EP");

        let result = server.graphql(
            Parameters(GraphqlParams {
                query: r#"{ tracks(search: "aphex") { edges { node { title artist } } } }"#.into(),
                variables: None,
            }),
            Default::default(),
        );
        assert!(result.is_ok());
        let Json(resp) = result.unwrap();
        let data = &resp.result["data"]["tracks"]["edges"];
        assert_eq!(data.as_array().unwrap().len(), 1);
        assert_eq!(data[0]["node"]["title"], "Windowlicker");
    }

    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
    async fn graphql_mutation_works() {
        let (server, _ch, _tmp) = test_server();
        let result = server.graphql(
            Parameters(GraphqlParams {
                query: "mutation { pause { ok message } }".into(),
                variables: None,
            }),
            Default::default(),
        );
        assert!(result.is_ok());
        let Json(resp) = result.unwrap();
        assert_eq!(resp.result["data"]["pause"]["ok"], true);
    }

    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
    async fn graphql_now_playing_stopped() {
        let (server, _ch, _tmp) = test_server();
        let result = server.graphql(
            Parameters(GraphqlParams {
                query: "{ nowPlaying { state positionMs } }".into(),
                variables: None,
            }),
            Default::default(),
        );
        assert!(result.is_ok());
        let Json(resp) = result.unwrap();
        assert_eq!(resp.result["data"]["nowPlaying"]["state"], "STOPPED");
    }

    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
    async fn graphql_library_stats() {
        let (server, _ch, tmp) = test_server();
        let db_path = tmp.path().join("test.db");
        insert_test_track(&db_path, "T1", "A1", "Album1");

        let result = server.graphql(
            Parameters(GraphqlParams {
                query: "{ libraryStats { totalTracks totalArtists totalAlbums } }".into(),
                variables: None,
            }),
            Default::default(),
        );
        assert!(result.is_ok());
        let Json(resp) = result.unwrap();
        assert_eq!(resp.result["data"]["libraryStats"]["totalTracks"], 1);
    }
}