1use std::path::PathBuf;
6use std::sync::Arc;
7
8use crate::auth::AuthUser;
9use crate::auth::password::Refused;
10use crossbeam_channel::Sender;
11use koan_core::player::commands::PlayerCommand;
12use koan_core::player::state::SharedPlayerState;
13use rmcp::handler::server::router::tool::ToolRouter;
14use rmcp::handler::server::wrapper::Json;
15use rmcp::model::{ServerCapabilities, ServerConfig};
16use rmcp::{ServerHandler, schemars, tool_router};
17use serde::{Deserialize, Serialize};
18
19#[derive(Debug, Deserialize, schemars::JsonSchema)]
24pub struct GraphqlParams {
25 #[schemars(
26 description = "GraphQL query or mutation string. Use the schema_sdl tool first to learn available types, queries, mutations, and filter parameters."
27 )]
28 pub query: String,
29 #[schemars(description = "Optional JSON object of query variables")]
30 pub variables: Option<serde_json::Value>,
31}
32
33#[derive(Debug, Serialize, schemars::JsonSchema)]
39pub struct GraphqlResponse {
40 pub result: serde_json::Value,
42}
43
44pub const USERNAME_HEADER: &str = "x-koan-username";
52pub const PASSWORD_HEADER: &str = "x-koan-password";
53
54#[derive(Clone)]
55pub struct KoanMcpServer {
56 #[allow(dead_code)]
57 tool_router: ToolRouter<Self>,
58 graphql_schema: crate::graphql::KoanSchema,
59 users: Option<Arc<crate::auth::password::PasswordVerifier>>,
61}
62
63impl KoanMcpServer {
64 pub fn new(
65 state: Arc<SharedPlayerState>,
66 cmd_tx: Sender<PlayerCommand>,
67 db_path: PathBuf,
68 ) -> Self {
69 let graphql_schema =
70 crate::graphql::build_schema(state.clone(), cmd_tx.clone(), db_path.clone(), None);
71 Self {
72 tool_router: Self::tool_router(),
73 graphql_schema,
74 users: None,
75 }
76 }
77
78 fn caller(&self, extensions: &rmcp::model::Extensions) -> Result<AuthUser, String> {
82 let local = AuthUser {
83 user_id: koan_core::db::queries::LOCAL_USER,
84 role: mcp_role(),
85 ..AuthUser::anonymous_admin()
86 };
87 let Some(users) = &self.users else {
88 return Ok(local);
89 };
90 let parts = extensions.get::<axum::http::request::Parts>();
91 let get = |h: &str| {
92 parts
93 .and_then(|p| p.headers.get(h))
94 .and_then(|v| v.to_str().ok())
95 .map(str::trim)
96 .filter(|v| !v.is_empty())
97 };
98 match (get(USERNAME_HEADER), get(PASSWORD_HEADER)) {
99 (Some(u), Some(p)) => match users.verify(u, p) {
100 Ok((user_id, role)) => Ok(AuthUser {
102 user_id,
103 username: u.to_owned(),
104 role,
105 }),
106 Err(Refused::Wrong) => Err("kōan rejected that username and password".into()),
107 Err(Refused::Busy) => Err("kōan is busy checking passwords; try again".into()),
108 },
109 _ if std::env::var("KOAN_MCP_REQUIRE_LOGIN").is_ok_and(|v| v == "1") => Err(format!(
110 "this kōan needs an account: send {USERNAME_HEADER} and {PASSWORD_HEADER}"
111 )),
112 _ => Ok(local),
113 }
114 }
115}
116
117use rmcp::handler::server::wrapper::Parameters;
118use rmcp::tool;
119
120fn mcp_role() -> koan_core::auth::Role {
128 if std::env::var("KOAN_MCP_ADMIN").is_ok_and(|v| v == "1") {
129 koan_core::auth::Role::Admin
130 } else {
131 koan_core::auth::Role::User
132 }
133}
134
135#[tool_router]
136impl KoanMcpServer {
137 #[tool(
138 description = "The GraphQL schema for the user's music (kōan): their library and the \
139 players they listen on. Call this first, before `graphql`. It covers playing, pausing, \
140 skipping and queueing music on the user's phone and computers, what is playing now, \
141 and searching, browsing and making playlists from the music they own."
142 )]
143 fn schema_sdl(&self) -> Json<GraphqlResponse> {
144 let sdl = self.graphql_schema.sdl();
145 Json(GraphqlResponse {
146 result: serde_json::Value::String(sdl),
147 })
148 }
149
150 #[tool(
151 description = "Control the user's music and search their music library (kōan). Use it \
152 for any request about music they listen to or own: play something, pause, resume, skip, \
153 what's playing, what's next, add to or change the queue, find or recommend from their \
154 collection, playlists, favourites. \"Pause the music on my desktop\", \"play some \
155 jazz on my phone\" and \"what is this song\" are all this tool.\n\n\
156 Call schema_sdl first for the full schema. The user's phones and computers running \
157 kōan are `clients`; commands for them end in `OnClient`.\n\n\
158 Examples:\n\
159 - What's playing, where: { clients { name playing nowPlaying positionMs } }\n\
160 - Pause: mutation { controlClient(action: PAUSE) { ok message } }\n\
161 - Find music: { tracks(search: \"aphex\", first: 20) { edges { node { id title artist album } } } }\n\
162 - Play it: mutation { playOnClient(trackIds: [\"42\", \"43\"]) { ok message } }\n\n\
163 String filters are case-insensitive substrings."
164 )]
165 fn graphql(
166 &self,
167 Parameters(params): Parameters<GraphqlParams>,
168 extensions: rmcp::model::Extensions,
169 ) -> Result<Json<GraphqlResponse>, String> {
170 let schema = self.graphql_schema.clone();
171 let query = params.query;
172 let variables = params.variables;
173 let rt =
174 tokio::runtime::Handle::try_current().map_err(|_| "no tokio runtime".to_string())?;
175 let result = tokio::task::block_in_place(|| {
177 let caller = self.caller(&extensions)?;
178 Ok::<_, String>(rt.block_on(crate::graphql::execute_in_process(
179 &schema, &query, variables, caller,
180 )))
181 })?;
182 Ok(Json(GraphqlResponse { result }))
183 }
184}
185
186#[rmcp::tool_handler]
187impl ServerHandler for KoanMcpServer {
188 fn get_info(&self) -> ServerConfig {
189 let instructions = if self.users.is_some() {
193 SERVER_INSTRUCTIONS
194 } else {
195 LOCAL_INSTRUCTIONS
196 };
197 ServerConfig::new(ServerCapabilities::builder().enable_tools().build())
198 .with_server_info(rmcp::model::Implementation::new(
199 "koan",
200 env!("CARGO_PKG_VERSION"),
201 ))
202 .with_instructions(instructions)
203 }
204}
205
206const SERVER_INSTRUCTIONS: &str = "kōan is the user's music: their whole music library, and the \
207phones and computers they listen on. Use it for anything about music they are playing or own — \
208\"pause the music\", \"play something like Polar Bear on my phone\", \"what's this song\", \
209\"skip to the Phace remix\", \"add their new album when it's downloaded\". Call `schema_sdl` \
210once, then do everything through `graphql`.
211
212## Where the music plays
213The user listens in kōan apps on their devices, linked to this server. Query \
214`clients { name platform playing nowPlaying album positionMs durationMs radio queue { trackId \
215title artist current } }` to see each device, what it is playing and what it has queued. Every \
216command about the user's music goes to a device:
217- `controlClient(action: PAUSE|RESUME|NEXT|PREVIOUS)`, `seekOnClient(positionMs)`
218- `playOnClient(trackIds, startAt)` replaces the queue and plays; `enqueue: true` appends. \
219A phone iOS has suspended is not linked but is still reached. Music comes up there as a \
220notification to tap, since iOS lets no app start audio on its own from sleep; queue, radio and \
221other changes are applied as it wakes. The message says when a device was asleep: tell the user \
222to tap the notification
223- `playNextOnClient(trackIds)`, `jumpOnClient(trackId)` (skip to a track, queued or not), \
224`removeFromClient(trackIds)`, `clearClient`, `setClientRadio(enabled)`, `syncClient`
225- **Making a playlist the user asked for** (\"make me a cyberpunk playlist\"): research what \
226fits, find each track in the library, `createPlaylist` with those in order. For picks the \
227library lacks, fetch the album with slsk's `grab`, then `addToPlaylistWhenAdded(playlistId, \
228artist, album, titles)` to add the wanted tracks once it is imported. Tell the user what is \
229there now and what is on its way.
230- Playlists made or edited here (`createPlaylist`, `setPlaylistTracks`…) reach every device \
231by themselves: linked ones sync at once, others when next opened. `syncClients` does the same \
232on request.
233- `evictOnClients(trackIds)` makes every linked device drop its downloaded copies of those \
234tracks: when a track plays as noise or glitches, after the file on the server is replaced
235- `queueOnClientWhenAdded(artist, album)` queues an album once it reaches the library, e.g. \
236one being downloaded with slsk's `grab`; `clientOrders` lists those waiting
237Leave `client` out unless the user named a device (\"my phone\", \"the desktop\": match it \
238against `clients` names and platforms). Without it the server picks the device that is \
239playing, else the one played most recently; if it answers that it cannot tell, ask the user \
240which device.
241
242**Act on what the user asks; do not second-guess it from reported state.** \"Pause\", \
243\"skip\" and \"resume\" go straight to `controlClient`: the user can hear the device and you \
244cannot, and a report can be stale or, from an older app (`playing: null`), absent.
245
246**Never use the server's own player for the user's music.** `play`, `pause`, `resume`, \
247`next`, `previous`, `seek`, `nowPlaying`, `queue`, `addToQueue`, `replaceQueue`, \
248`playPlaylist` and the radio mutations drive a headless player on the server that nobody \
249hears; `nowPlaying` there reports nothing about what the user is listening to.
250
251## The library
252- `artists`, `albums`, `tracks` with filters (genre, year range, codec, sample rate, bit depth, \
253duration, favourites), `randomTracks`, `similarArtists`, `similarTracks`, `fuzzySearch`
254- Build a set from these, then send its track ids to a device with `playOnClient`. Track ids are \
255integers in queries; pass them to the client mutations as strings.
256- Favourites: `favourite`, `unfavourite`, `toggleFavourite`, `favouritesOnly: true` on queries
257- Playlists: `playlists`, `playlistTracks`, `createPlaylist`, `addToPlaylist`, \
258`setPlaylistTracks`, `renamePlaylist`, `deletePlaylist`
259- History: `playHistory`
260- Sharing: `createShare(trackIds, description)` makes a public link anyone can open without an \
261account; confirm with the user first. `shares`, `updateShare`, `deleteShare` manage them.
262
263## Not available
264`organize*` (moves files on disk), `updateConfig` and `triggerScan` are refused unless \
265`KOAN_MCP_ADMIN=1` is set.";
266
267const LOCAL_INSTRUCTIONS: &str = "kōan is the user's music player on this machine and their \
268music library. Use it for anything about music they are playing or own — \"pause the music\", \
269\"play something like Polar Bear\", \"what's this song\". Call `schema_sdl` once, then do \
270everything through `graphql`.
271
272## Playback
273This player is what the user hears: `play`, `pause`, `resume`, `stop`, `next`, `previous`, \
274`seek`, `nowPlaying`; the queue with `queue`, `addToQueue`, `replaceQueue`, `removeFromQueue`, \
275`moveInQueue`, `clearQueue`, `undo`, `redo`; radio with `enableRadio`, `disableRadio`.
276
277## The library
278- `artists`, `albums`, `tracks` with filters (genre, year range, codec, sample rate, bit depth, \
279duration, favourites), `randomTracks`, `similarArtists`, `similarTracks`, `fuzzySearch`
280- Favourites: `favourite`, `unfavourite`, `toggleFavourite`, `favouritesOnly: true` on queries
281- Playlists: `playlists`, `playlistTracks`, `createPlaylist`, `saveQueueAsPlaylist`, \
282`addToPlaylist`, `setPlaylistTracks`, `renamePlaylist`, `deletePlaylist`, `playPlaylist`
283- History: `playHistory`
284- Sharing: `createShare(trackIds, description)` makes a public link; confirm with the user first.
285
286## Not available
287`organize*` (moves files on disk), `updateConfig`, `triggerScan` and `setDevice` are refused \
288unless `KOAN_MCP_ADMIN=1` is set.
289
290## IDs
291Track IDs are integers from the library; queue item IDs are UUIDs from the queue.";
292
293pub fn spawn_http(
301 addr: std::net::SocketAddr,
302 state: Arc<SharedPlayerState>,
303 cmd_tx: Sender<PlayerCommand>,
304 pool: Arc<koan_core::db::pool::Pool>,
305) -> std::io::Result<std::thread::JoinHandle<()>> {
306 use rmcp::transport::streamable_http_server::{
307 StreamableHttpServerConfig, StreamableHttpService, session::local::LocalSessionManager,
308 };
309 let mut template = KoanMcpServer::new(state, cmd_tx, pool.path().to_path_buf());
310 template.users = Some(Arc::new(crate::auth::password::PasswordVerifier::new(pool)));
311 let listener = std::net::TcpListener::bind(addr)?;
313 listener.set_nonblocking(true)?;
314 std::thread::Builder::new()
315 .name("koan-mcp-http".into())
316 .spawn(move || {
317 let rt = tokio::runtime::Builder::new_multi_thread()
318 .enable_all()
319 .build()
320 .expect("failed to create tokio runtime");
321 rt.block_on(async move {
322 let service = StreamableHttpService::new(
323 move || Ok(template.clone()),
324 Arc::new(LocalSessionManager::default()),
325 StreamableHttpServerConfig::default().disable_allowed_hosts(),
329 );
330 let app = axum::Router::new().nest_service("/mcp", service);
331 let listener =
332 tokio::net::TcpListener::from_std(listener).expect("listener from std");
333 if let Err(e) = axum::serve(listener, app).await {
334 log::error!("MCP HTTP server stopped: {e}");
335 }
336 });
337 })
338}
339
340pub fn cmd_mcp() {
342 use koan_core::player::Player;
343 use rmcp::ServiceExt;
344
345 let _db = koan_core::db::connection::Database::open_default().expect("failed to open database");
347 let db_path = koan_core::config::db_path();
348
349 let (state, _timeline, _viz, cmd_tx) = Player::spawn();
351
352 let server = KoanMcpServer::new(state, cmd_tx, db_path);
353
354 let rt = tokio::runtime::Runtime::new().expect("failed to create tokio runtime");
356 rt.block_on(async {
357 let transport = rmcp::transport::io::stdio();
358 let service = server
359 .serve(transport)
360 .await
361 .expect("failed to start MCP server");
362 let _ = service.waiting().await;
363 });
364}
365
366#[cfg(test)]
371mod tests {
372 use super::*;
373 use koan_core::db::connection::Database;
374 use koan_core::db::queries;
375 use koan_core::player::commands::CommandChannel;
376 use tempfile::TempDir;
377
378 fn test_server() -> (KoanMcpServer, CommandChannel, TempDir) {
379 let tmp = TempDir::new().unwrap();
380 let db_path = tmp.path().join("test.db");
381 let db = Database::open(&db_path).unwrap();
382 koan_core::db::schema::create_tables(&db.conn).unwrap();
383
384 let state = SharedPlayerState::new();
385 let ch = CommandChannel::new();
386 let tx = ch.tx.clone();
387
388 let server = KoanMcpServer::new(state, tx, db_path);
389 (server, ch, tmp)
390 }
391
392 fn with_headers(headers: &[(&str, &str)]) -> rmcp::model::Extensions {
393 let mut req = axum::http::Request::builder();
394 for (k, v) in headers {
395 req = req.header(*k, *v);
396 }
397 let (parts, ()) = req.body(()).unwrap().into_parts();
398 let mut ext = rmcp::model::Extensions::new();
399 ext.insert(parts);
400 ext
401 }
402
403 #[test]
404 fn gateway_headers_act_as_that_account() {
405 use koan_core::auth::Role;
406 let (mut server, _ch, tmp) = test_server();
407 let db_path = tmp.path().join("test.db");
408 let db = Database::open(&db_path).unwrap();
409 queries::auth::create_user(&db.conn, "owner", "sesame", Role::Admin).unwrap();
410 queries::auth::create_user(&db.conn, "mate", "hunter22", Role::Readonly).unwrap();
411 server.users = Some(Arc::new(crate::auth::password::PasswordVerifier::new(
412 Arc::new(koan_core::db::pool::Pool::new(db_path)),
413 )));
414
415 let as_ = |u: &str, p: &str| {
416 server
417 .caller(&with_headers(&[(USERNAME_HEADER, u), (PASSWORD_HEADER, p)]))
418 .map(|c| (c.user_id, c.username, c.role))
419 };
420 assert_eq!(as_("owner", "sesame"), Ok((1, "owner".into(), Role::Admin)));
422 assert_eq!(
423 as_("mate", "hunter22"),
424 Ok((2, "mate".into(), Role::Readonly))
425 );
426 assert!(as_("owner", "wrong").is_err());
427 let local = server.caller(&with_headers(&[])).unwrap();
429 assert_eq!(
430 (local.user_id, local.role),
431 (queries::LOCAL_USER, mcp_role())
432 );
433 }
434
435 fn insert_test_track(db_path: &std::path::Path, title: &str, artist: &str, album: &str) -> i64 {
436 let db = Database::open(db_path).unwrap();
437 let meta = queries::TrackMeta {
438 title: title.to_string(),
439 artist: artist.to_string(),
440 album_artist: Some(artist.to_string()),
441 album: album.to_string(),
442 track_number: Some(1),
443 disc: Some(1),
444 date: Some("2024".into()),
445 genre: Some("Electronic".into()),
446 duration_ms: Some(240000),
447 path: Some(format!(
448 "/tmp/test/{}.flac",
449 title.to_lowercase().replace(' ', "_")
450 )),
451 codec: Some("FLAC".into()),
452 sample_rate: Some(44100),
453 bit_depth: Some(16),
454 channels: Some(2),
455 bitrate: Some(1411),
456 size_bytes: Some(42_000_000),
457 mtime: Some(1700000000),
458 source: "local".into(),
459 remote_id: None,
460 remote_url: None,
461 album_remote_id: None,
462 artist_remote_id: None,
463 mbid: None,
464 album_mbid: None,
465 album_added_at: None,
466 label: None,
467 };
468 queries::upsert_track(&db.conn, &meta).unwrap()
469 }
470
471 #[test]
472 fn schema_sdl_returns_schema() {
473 let (server, _ch, _tmp) = test_server();
474 let Json(resp) = server.schema_sdl();
475 let sdl = resp.result.as_str().unwrap();
476 assert!(sdl.contains("type QueryRoot"));
477 assert!(sdl.contains("type MutationRoot"));
478 assert!(sdl.contains("artists"));
479 assert!(sdl.contains("nowPlaying"));
480 }
481
482 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
483 async fn graphql_query_works() {
484 let (server, _ch, tmp) = test_server();
485 let db_path = tmp.path().join("test.db");
486 insert_test_track(&db_path, "Windowlicker", "Aphex Twin", "Windowlicker EP");
487
488 let result = server.graphql(
489 Parameters(GraphqlParams {
490 query: r#"{ tracks(search: "aphex") { edges { node { title artist } } } }"#.into(),
491 variables: None,
492 }),
493 Default::default(),
494 );
495 assert!(result.is_ok());
496 let Json(resp) = result.unwrap();
497 let data = &resp.result["data"]["tracks"]["edges"];
498 assert_eq!(data.as_array().unwrap().len(), 1);
499 assert_eq!(data[0]["node"]["title"], "Windowlicker");
500 }
501
502 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
503 async fn graphql_mutation_works() {
504 let (server, _ch, _tmp) = test_server();
505 let result = server.graphql(
506 Parameters(GraphqlParams {
507 query: "mutation { pause { ok message } }".into(),
508 variables: None,
509 }),
510 Default::default(),
511 );
512 assert!(result.is_ok());
513 let Json(resp) = result.unwrap();
514 assert_eq!(resp.result["data"]["pause"]["ok"], true);
515 }
516
517 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
518 async fn graphql_now_playing_stopped() {
519 let (server, _ch, _tmp) = test_server();
520 let result = server.graphql(
521 Parameters(GraphqlParams {
522 query: "{ nowPlaying { state positionMs } }".into(),
523 variables: None,
524 }),
525 Default::default(),
526 );
527 assert!(result.is_ok());
528 let Json(resp) = result.unwrap();
529 assert_eq!(resp.result["data"]["nowPlaying"]["state"], "STOPPED");
530 }
531
532 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
533 async fn graphql_library_stats() {
534 let (server, _ch, tmp) = test_server();
535 let db_path = tmp.path().join("test.db");
536 insert_test_track(&db_path, "T1", "A1", "Album1");
537
538 let result = server.graphql(
539 Parameters(GraphqlParams {
540 query: "{ libraryStats { totalTracks totalArtists totalAlbums } }".into(),
541 variables: None,
542 }),
543 Default::default(),
544 );
545 assert!(result.is_ok());
546 let Json(resp) = result.unwrap();
547 assert_eq!(resp.result["data"]["libraryStats"]["totalTracks"], 1);
548 }
549}