1//! An MCP server inside the window, so that whoever is developing the app can 2//! ask the running console what it is doing instead of photographing it. 3//! 4//! The console belongs to the UI thread and stays there. A tool call becomes a 5//! [`Request`] on a channel, the UI thread answers it at a frame boundary, and 6//! the reply comes back on a oneshot. Nothing here touches the emulator. 7//! 8//! Who is developing the app is not who plays the game. The app's own loops 9//! play; a client watches. So the tools come in two sets: the ones that only 10//! look are always there, and the ones that change anything — the pad, the 11//! machine's state, the process — exist only while [`DevMode`] is on, which it 12//! is not until somebody turns it on. Turning it either way tells every client 13//! that its tool list has changed. 14//! 15//! How a client is told depends on which MCP it speaks, and the transport is 16//! streamable HTTP either way. Before 2026-07-28 a client holds a stream open 17//! on its session, so there are sessions, kept in a file ([`crate::sessions`]) 18//! so that a client connected before the app restarted is still connected 19//! after it. From 2026-07-28 there are no sessions and a client asks to be told 20//! with `subscriptions/listen`, which it makes again itself after a restart. 21//! 22//! This crate is the whole native-only side of talking to the app from 23//! outside it: the server, `native mcp`'s stdio proxy ([`proxy`]) and the 24//! session/snapshot-name files ([`states`]). It is deliberately its OWN 25//! buck2 target, separate from `apps/native`'s tip crate — see this 26//! directory's CLAUDE.md for why that split exists and what it costs the 27//! hot-patch cycle to get wrong. 28 29pub mod proxy; 30mod sessions; 31pub mod states; 32 33use std::collections::VecDeque; 34use std::num::NonZeroU32; 35use std::sync::atomic::{AtomicBool, Ordering}; 36use std::sync::{Arc, Mutex, MutexGuard, OnceLock, PoisonError, mpsc}; 37use std::thread; 38use std::time::Instant; 39 40use base64::Engine; 41use console::SnesButton; 42use panels::pad::Confidence; 43use rmcp::handler::server::router::tool::ToolRouter; 44use rmcp::handler::server::wrapper::Parameters; 45use rmcp::handler::server::tool::ToolCallContext; 46use rmcp::model::{ 47 CallToolRequestParams, CallToolResponse, CallToolResult, ContentBlock, Implementation, 48 InitializeRequestParams, InitializeResult, ListToolsResult, PaginatedRequestParams, 49 ServerCapabilities, ServerConfig, SubscriptionFilter, Tool, 50}; 51use rmcp::service::{Peer, RequestContext, SubscriptionContext, SubscriptionSink}; 52use rmcp::transport::streamable_http_server::session::local::LocalSessionManager; 53use rmcp::transport::streamable_http_server::{StreamableHttpServerConfig, StreamableHttpService}; 54use rmcp::{ErrorData, RoleServer, ServerHandler, schemars, tool, tool_router}; 55use tokio::sync::oneshot; 56 57use crate::sessions::FileSessionStore; 58 59/// Where the server listens. Loopback only: this drives the pad and reads 60/// memory, and is for the machine the window is on. 61pub const ADDRESS: &str = "127.0.0.1:7637"; 62 63/// A request here is never answered, only dropped when the window goes. 64pub const ALIVE: &str = "/alive"; 65 66/// The longest a single `press` may run, in frames: one minute of game time. 67/// A tool call that outlives its client's patience reports nothing at all. 68const LONGEST_PRESS: u32 = 3600; 69 70/// Work RAM as the SNES addresses it. 71const WRAM_BUS: std::ops::Range<u32> = 0x7E0000..0x800000; 72 73/// How many tool calls the window keeps to show. 74const CALLS_KEPT: usize = 200; 75 76/// Whether the server is there to be talked to. 77#[derive(Debug, Default)] 78pub enum Listening { 79 #[default] 80 Starting, 81 On, 82 /// Why not: the port was taken, most likely by another window. 83 Failed(String), 84} 85 86#[derive(Debug)] 87pub enum Outcome { 88 Running, 89 Done, 90 Failed(String), 91} 92 93/// One tool call, as it arrived. 94#[derive(Debug)] 95pub struct Call { 96 id: u64, 97 pub at: Instant, 98 pub tool: String, 99 /// The arguments as the JSON they came in as; empty when there were none. 100 pub arguments: String, 101 pub outcome: Outcome, 102} 103 104/// Whether the tools that change things exist. Off when the window opens, 105/// whatever it was when the last one closed. 106#[derive(Debug, Default)] 107pub struct DevMode(AtomicBool); 108 109impl DevMode { 110 pub fn is_on(&self) -> bool { 111 self.0.load(Ordering::Relaxed) 112 } 113} 114 115/// What the window shows about the server. The server thread writes it and 116/// the UI thread reads it. 117#[derive(Debug, Default)] 118pub struct Activity { 119 pub listening: Listening, 120 pub dev_mode: Arc<DevMode>, 121 /// The client that last introduced itself, by name and version. A client 122 /// that was already connected before this window started never does. 123 pub client: Option<String>, 124 pub last_heard: Option<Instant>, 125 /// Oldest first, the newest [`CALLS_KEPT`]. 126 pub calls: VecDeque<Call>, 127 next_id: u64, 128} 129 130/// A panic elsewhere while the record was held leaves a record that is still 131/// worth showing, so poisoning is not an error here. 132pub fn record(activity: &Mutex<Activity>) -> MutexGuard<'_, Activity> { 133 activity.lock().unwrap_or_else(PoisonError::into_inner) 134} 135 136/// An RGBA picture, rows top to bottom. 137pub struct Picture { 138 pub width: u32, 139 pub height: u32, 140 pub rgba: Vec<u8>, 141} 142 143/// What the server asks of the UI thread. 144pub enum Request { 145 State(oneshot::Sender<alttp::State>), 146 /// The game's own picture, at the console's resolution. 147 Frame(oneshot::Sender<Picture>), 148 /// The whole window as drawn, panels included. 149 Window(oneshot::Sender<Picture>), 150 Wram { 151 /// Offset into work RAM; in range by construction. 152 at: usize, 153 len: usize, 154 reply: oneshot::Sender<Vec<u8>>, 155 }, 156 /// The names of the snapshots there are. 157 States(oneshot::Sender<Vec<String>>), 158 /// What the zbanks bot says about itself (`zbanks::Report`), already 159 /// serialized to JSON by the app. Read-only. 160 Bot(oneshot::Sender<String>), 161 /// Jev's recent goal-choice history (`decisions::goal_choice::Chooser::history`), 162 /// already serialized to JSON by the app, newest first. Read-only. 163 JevHistory(oneshot::Sender<String>), 164 /// Hold `buttons` for `frames` emulated frames, release them, then reply 165 /// with the state the game is left in. 166 Press { 167 buttons: Vec<SnesButton>, 168 frames: NonZeroU32, 169 confidence: Option<Confidence>, 170 reply: oneshot::Sender<alttp::State>, 171 }, 172 /// Keep the machine as it stands under a name. 173 SaveState { name: states::Name, reply: oneshot::Sender<Result<(), String>> }, 174 /// Become the machine kept under a name, and reply with the game's state. 175 LoadState { name: states::Name, reply: oneshot::Sender<Result<alttp::State, String>> }, 176 /// Turn the machine off and on again. 177 PowerCycle(oneshot::Sender<alttp::State>), 178 /// Keep where the game is, then replace this process with a fresh run of 179 /// the binary it was started from, which carries on from there. 180 Restart(oneshot::Sender<Result<(), String>>), 181 /// Apply a subsecond hot patch. Must run on the UI thread at a frame 182 /// boundary, like every other request here - `subsecond::apply_patch` 183 /// rewrites which code the running `App` value's patch points jump to, 184 /// and that jump table is read from the same thread that calls through 185 /// it (`apps/native/src/app.rs`'s `logic_impl`/`ui_impl`). 186 HotPatch { table: subsecond_types::JumpTable, reply: oneshot::Sender<Result<usize, String>> }, 187} 188 189#[derive(Debug, Clone, Copy, serde::Deserialize, schemars::JsonSchema)] 190#[serde(rename_all = "lowercase")] 191enum Button { 192 Up, 193 Down, 194 Left, 195 Right, 196 A, 197 B, 198 X, 199 Y, 200 L, 201 R, 202 Start, 203 Select, 204} 205 206impl From<Button> for SnesButton { 207 fn from(button: Button) -> Self { 208 match button { 209 Button::Up => Self::Up, 210 Button::Down => Self::Down, 211 Button::Left => Self::Left, 212 Button::Right => Self::Right, 213 Button::A => Self::A, 214 Button::B => Self::B, 215 Button::X => Self::X, 216 Button::Y => Self::Y, 217 Button::L => Self::L, 218 Button::R => Self::R, 219 Button::Start => Self::Start, 220 Button::Select => Self::Select, 221 } 222 } 223} 224 225#[derive(Debug, serde::Deserialize, schemars::JsonSchema)] 226struct PressRequest { 227 #[schemars(description = "Buttons to hold together. Empty lets the game run with the pad idle.")] 228 buttons: Vec<Button>, 229 #[schemars(description = "Emulated frames to hold them for; the game runs at 60 a second.")] 230 frames: NonZeroU32, 231 #[schemars(description = "How sure you are of this press, as a probability from 0 to 1. The window shows it on the controller.")] 232 confidence: Option<f32>, 233} 234 235#[derive(Debug, serde::Deserialize, schemars::JsonSchema)] 236struct WramRequest { 237 #[schemars(description = "SNES bus address in hex, as the RAM map writes it: \"7EF36C\" or \"$7EF36C\".")] 238 address: String, 239 #[schemars(description = "How many bytes to read.")] 240 length: u16, 241} 242 243#[derive(Debug, serde::Deserialize, schemars::JsonSchema)] 244struct DevModeRequest { 245 #[schemars(description = "true adds the tools that change things; false takes them away.")] 246 on: bool, 247} 248 249#[derive(Debug, serde::Deserialize, schemars::JsonSchema)] 250struct StateName { 251 #[schemars(description = "Lower-case letters, digits, '-' and '_'. \"resume\" is the one the app keeps for itself and starts from.")] 252 name: String, 253} 254 255#[derive(Debug, serde::Deserialize, schemars::JsonSchema)] 256struct HotPatchRequest { 257 #[schemars(description = "Path to a JSON-serialized subsecond_types::JumpTable, as tools/hotpatch writes it (patch-N.json).")] 258 jump_table: String, 259} 260 261#[derive(Debug, serde::Serialize, schemars::JsonSchema)] 262struct PatchInfo { 263 pid: u32, 264 /// `subsecond::aslr_reference()` in THIS process: `main`'s live address, 265 /// which a patch build rebases its jump table against 266 /// (research/subsecond-patch-build.md §2/§7). 267 aslr_reference: String, 268} 269 270/// Somebody to tell that the tool list changed. 271enum Listener { 272 /// A session, which lasts until telling it fails. 273 Session(Peer<RoleServer>), 274 /// A `subscriptions/listen`, which lasts until its request is cancelled. 275 Subscription { id: u64, sink: SubscriptionSink }, 276} 277 278impl Listener { 279 async fn tell(&self) -> bool { 280 match self { 281 Self::Session(peer) => peer.notify_tool_list_changed().await.is_ok(), 282 Self::Subscription { sink, .. } => sink.notify_tool_list_changed().await.is_ok(), 283 } 284 } 285} 286 287#[derive(Default)] 288struct Listeners { 289 all: Vec<Listener>, 290 next_subscription: u64, 291} 292 293type Peers = Arc<Mutex<Listeners>>; 294 295/// One per session. 296#[derive(Clone)] 297struct Server { 298 requests: mpsc::Sender<Request>, 299 ui: egui::Context, 300 activity: Arc<Mutex<Activity>>, 301 dev_mode: Arc<DevMode>, 302 peers: Peers, 303 /// Set once this session's peer is among `peers`. 304 known: Arc<OnceLock<()>>, 305 watching: Arc<ToolRouter<Self>>, 306 driving: Arc<ToolRouter<Self>>, 307} 308 309impl Server { 310 /// The router a tool is in, if it exists right now. A driving tool called 311 /// while dev mode is off is not refused, it is not there. 312 fn router_of(&self, tool: &str) -> Option<&ToolRouter<Self>> { 313 if self.watching.has_route(tool) { 314 Some(&self.watching) 315 } else if self.dev_mode.is_on() && self.driving.has_route(tool) { 316 Some(&self.driving) 317 } else { 318 None 319 } 320 } 321 322 /// Any request carries the peer to reach its session by, including the 323 /// first one from a session restored after a restart, which never 324 /// introduces itself again. 325 fn meet(&self, peer: &Peer<RoleServer>) { 326 self.known.get_or_init(|| { 327 let mut peers = self.peers.lock().unwrap_or_else(PoisonError::into_inner); 328 peers.all.push(Listener::Session(peer.clone())); 329 }); 330 } 331 332 fn state_name(name: &str) -> Result<states::Name, ErrorData> { 333 states::Name::new(name).ok_or_else(|| { 334 ErrorData::invalid_params( 335 format!("{name:?} is not a state name: lower-case letters, digits, '-' and '_'"), 336 None, 337 ) 338 }) 339 } 340 341 /// Note that a client was heard from, record what it said, and have the 342 /// window redraw to show it. 343 fn heard<T>(&self, note: impl FnOnce(&mut Activity) -> T) -> T { 344 let noted = { 345 let mut activity = record(&self.activity); 346 activity.last_heard = Some(Instant::now()); 347 note(&mut activity) 348 }; 349 self.ui.request_repaint(); 350 noted 351 } 352 353 /// Put a request to the UI thread and wait for its answer. 354 async fn ask<T>( 355 &self, 356 request: impl FnOnce(oneshot::Sender<T>) -> Request, 357 ) -> Result<T, ErrorData> { 358 fn gone<E>(_: E) -> ErrorData { 359 ErrorData::internal_error("the window has closed", None) 360 } 361 let (reply, answer) = oneshot::channel(); 362 self.requests.send(request(reply)).map_err(gone)?; 363 // The UI thread sleeps between frames; this is what it is to wake for. 364 self.ui.request_repaint(); 365 answer.await.map_err(gone) 366 } 367} 368 369fn png(picture: &Picture) -> Result<CallToolResult, ErrorData> { 370 let failed = |e: png::EncodingError| ErrorData::internal_error(e.to_string(), None); 371 let mut bytes = Vec::new(); 372 let mut encoder = png::Encoder::new(&mut bytes, picture.width, picture.height); 373 encoder.set_color(png::ColorType::Rgba); 374 encoder.set_depth(png::BitDepth::Eight); 375 let mut writer = encoder.write_header().map_err(failed)?; 376 writer.write_image_data(&picture.rgba).map_err(failed)?; 377 writer.finish().map_err(failed)?; 378 let encoded = base64::engine::general_purpose::STANDARD.encode(bytes); 379 Ok(CallToolResult::success(vec![ContentBlock::image(encoded, "image/png")])) 380} 381 382#[tool_router(router = watching)] 383impl Server { 384 #[tool(description = "What the game is doing and where Link is, read from work RAM as of the last frame.")] 385 async fn state(&self) -> Result<CallToolResult, ErrorData> { 386 let state = self.ask(Request::State).await?; 387 Ok(CallToolResult::success(vec![ContentBlock::json(state)?])) 388 } 389 390 #[tool(description = "The game's picture as of the last frame, as a PNG at the console's own resolution.")] 391 async fn frame(&self) -> Result<CallToolResult, ErrorData> { 392 png(&self.ask(Request::Frame).await?) 393 } 394 395 #[tool(description = "The whole app window as it is drawn, panels included, as a PNG.")] 396 async fn window(&self) -> Result<CallToolResult, ErrorData> { 397 png(&self.ask(Request::Window).await?) 398 } 399 400 #[tool(description = "Read bytes of work RAM ($7E0000-$7FFFFF), returned as hex.")] 401 async fn read_wram( 402 &self, 403 Parameters(WramRequest { address, length }): Parameters<WramRequest>, 404 ) -> Result<CallToolResult, ErrorData> { 405 let bus = u32::from_str_radix(address.trim_start_matches('$'), 16).map_err(|e| { 406 ErrorData::invalid_params(format!("address {address:?} is not hex: {e}"), None) 407 })?; 408 // Checked in this order so the sum below cannot overflow. 409 let end = WRAM_BUS.contains(&bus).then(|| bus + u32::from(length)); 410 if !end.is_some_and(|end| end <= WRAM_BUS.end) { 411 return Err(ErrorData::invalid_params( 412 format!("{length} bytes at ${bus:06X} is not inside work RAM, $7E0000..$800000"), 413 None, 414 )); 415 } 416 let at = (bus - WRAM_BUS.start) as usize; 417 let bytes = self 418 .ask(|reply| Request::Wram { at, len: usize::from(length), reply }) 419 .await?; 420 let hex: Vec<String> = bytes.iter().map(|b| format!("{b:02X}")).collect(); 421 Ok(CallToolResult::success(vec![ContentBlock::text(hex.join(" "))])) 422 } 423 424 #[tool(description = "The names of the machine snapshots kept beside the ROM.")] 425 async fn states(&self) -> Result<CallToolResult, ErrorData> { 426 let names = self.ask(Request::States).await?; 427 Ok(CallToolResult::success(vec![ContentBlock::json(names)?])) 428 } 429 430 #[tool(description = "What the bot playing the game (zbanks/alttp, C, linked in) says about itself, as JSON: its info line, its task list (current first), the first goals on its goal list with type, node, screen, attempts and last score, how many goals there are, Link where the bot places him (its own map coordinates), how many assert_bp traps have fired, and base() reads that mapped to nothing. Read-only.")] 431 async fn bot(&self) -> Result<CallToolResult, ErrorData> { 432 let json = self.ask(Request::Bot).await?; 433 Ok(CallToolResult::success(vec![ContentBlock::text(json)])) 434 } 435 436 #[tool(description = "Jev's recent goal-choice history, newest first, bounded to the same count the Bot panel's Jev history shows (older decisions stay in the run's own JSONL log on disk, not here): for each choice, the frame, every option in the words Jev was given (upstream's own pick first), Jev's probability per option when it asked, the one picked, and how the choice was made - asked (with the exact request sent, tokens, latency and cost), reused from a recent answer, one option after look-alikes collapsed, throttled, or upstream's own pick with why. Read-only.")] 437 async fn jev_history(&self) -> Result<CallToolResult, ErrorData> { 438 let json = self.ask(Request::JevHistory).await?; 439 Ok(CallToolResult::success(vec![ContentBlock::text(json)])) 440 } 441 442 #[tool(description = "Jev's spend guards, read from the shared ledger every process writes to (so every process's questions count, not only this window's): spend so far, plainly, with no cap or denominator (there is no self-imposed lifetime cap - the only hard stop is the vendor's own out-of-credit answer); the question bucket (questions in the last minute, its size, tokens available, and when empty how many seconds until the next token returns - it refills by itself as the minute slides); the hour breaker (dollars in the last hour, its limit, and when closed how many seconds until it reopens by itself); holds in flight; whether the vendor said the account is out of credit; the last rate-limit snapshot; the ledger's path; and one line saying what, if anything, is stopping questions. Read-only: nothing here needs clearing by hand.")] 443 async fn budget(&self) -> Result<CallToolResult, ErrorData> { 444 let opened = |e: String| ErrorData::internal_error(format!("the spend ledger: {e}"), None); 445 let ledger = jev_http::Ledger::open().map_err(opened)?; 446 let status = ledger 447 .status(jev_http::ledger::now(), jev_http::Guards::from_env()) 448 .map_err(opened)?; 449 // Flat, so tools/watch/budget.sh reads top-level fields. 450 let mut budget = serde_json::to_value(&status) 451 .map_err(|e| ErrorData::internal_error(format!("the guards' status: {e}"), None))?; 452 if let Some(fields) = budget.as_object_mut() { 453 fields.insert("stopped".into(), serde_json::json!(status.stopped())); 454 fields.insert("stopped_for_good".into(), serde_json::json!(status.stopped_for_good())); 455 fields.insert("line".into(), serde_json::json!(status.line())); 456 fields.insert("last_rate_limit".into(), serde_json::json!(ledger.last_rate_limit())); 457 fields.insert("ledger_path".into(), serde_json::json!(ledger.path().display().to_string())); 458 } 459 Ok(CallToolResult::success(vec![ContentBlock::json(budget)?])) 460 } 461 462 #[tool(description = "This process's pid and subsecond::aslr_reference(), for tools/hotpatch/patch.sh to build a patch against - no log scraping needed.")] 463 async fn patch_info(&self) -> Result<CallToolResult, ErrorData> { 464 let info = PatchInfo { 465 pid: std::process::id(), 466 aslr_reference: format!("{:#x}", subsecond::aslr_reference()), 467 }; 468 Ok(CallToolResult::success(vec![ContentBlock::json(info)?])) 469 } 470 471 #[tool(description = "Dev mode is off when the app starts: the app plays the game and these tools only watch. Turning it on adds the tools that change things (press, save_state, load_state, power_cycle, restart), for debugging the app while building it.")] 472 async fn dev_mode( 473 &self, 474 Parameters(DevModeRequest { on }): Parameters<DevModeRequest>, 475 ) -> Result<CallToolResult, ErrorData> { 476 if self.dev_mode.0.swap(on, Ordering::Relaxed) != on { 477 // Taken out so the lock is not held while sending. One that 478 // cannot be told has gone away, and is forgotten. 479 let lock = || self.peers.lock().unwrap_or_else(PoisonError::into_inner); 480 let listeners = std::mem::take(&mut lock().all); 481 let mut reached = Vec::with_capacity(listeners.len()); 482 for listener in listeners { 483 if listener.tell().await { 484 reached.push(listener); 485 } 486 } 487 lock().all.append(&mut reached); 488 self.ui.request_repaint(); 489 } 490 let now = if on { "on" } else { "off" }; 491 Ok(CallToolResult::success(vec![ContentBlock::text(format!("dev mode is {now}"))])) 492 } 493} 494 495#[tool_router(router = driving)] 496impl Server { 497 #[tool(description = "Hold buttons on controller one for a number of frames, release them, and return the state the game is left in. With no buttons it just lets the game run.")] 498 async fn press( 499 &self, 500 Parameters(PressRequest { buttons, frames, confidence }): Parameters<PressRequest>, 501 ) -> Result<CallToolResult, ErrorData> { 502 let confidence = confidence 503 .map(|probability| { 504 Confidence::new(probability).ok_or_else(|| { 505 ErrorData::invalid_params( 506 format!("confidence is {probability}; a probability is 0 to 1"), 507 None, 508 ) 509 }) 510 }) 511 .transpose()?; 512 if frames.get() > LONGEST_PRESS { 513 return Err(ErrorData::invalid_params( 514 format!("frames is {frames}; the most one call may hold is {LONGEST_PRESS}"), 515 None, 516 )); 517 } 518 let buttons = buttons.into_iter().map(SnesButton::from).collect(); 519 let state = self 520 .ask(|reply| Request::Press { buttons, frames, confidence, reply }) 521 .await?; 522 Ok(CallToolResult::success(vec![ContentBlock::json(state)?])) 523 } 524 525 #[tool(description = "Keep the machine exactly as it stands, under a name, beside the ROM.")] 526 async fn save_state( 527 &self, 528 Parameters(StateName { name }): Parameters<StateName>, 529 ) -> Result<CallToolResult, ErrorData> { 530 let name = Self::state_name(&name)?; 531 self.ask(|reply| Request::SaveState { name, reply }) 532 .await? 533 .map_err(|e| ErrorData::internal_error(e, None))?; 534 Ok(CallToolResult::success(vec![ContentBlock::text("kept")])) 535 } 536 537 #[tool(description = "Become the machine kept under a name, and return the state the game is then in.")] 538 async fn load_state( 539 &self, 540 Parameters(StateName { name }): Parameters<StateName>, 541 ) -> Result<CallToolResult, ErrorData> { 542 let name = Self::state_name(&name)?; 543 let state = self 544 .ask(|reply| Request::LoadState { name, reply }) 545 .await? 546 .map_err(|e| ErrorData::invalid_params(e, None))?; 547 Ok(CallToolResult::success(vec![ContentBlock::json(state)?])) 548 } 549 550 #[tool(description = "Turn the SNES off and on again. The cartridge's battery save is kept; where the game was is not, unless save_state kept it first.")] 551 async fn power_cycle(&self) -> Result<CallToolResult, ErrorData> { 552 let state = self.ask(Request::PowerCycle).await?; 553 Ok(CallToolResult::success(vec![ContentBlock::json(state)?])) 554 } 555 556 #[tool(description = "Restart the app on the binary it was started from, to run a new build: the game carries on from the frame it was on, this session stays connected, and dev mode is off again.")] 557 async fn restart(&self) -> Result<CallToolResult, ErrorData> { 558 self.ask(Request::Restart).await?.map_err(|e| ErrorData::internal_error(e, None))?; 559 Ok(CallToolResult::success(vec![ContentBlock::text("restarting")])) 560 } 561 562 #[tool(description = "Apply a subsecond hot patch (tools/hotpatch's patch-N.json) to the running window: no restart, the game keeps running. Only rewrites function bodies behind a subsecond::call/HotFn boundary (apps/native/src/app.rs's logic_impl/ui_impl, packages/panels' patch points) - a struct-layout change needs restart instead.")] 563 async fn hot_patch( 564 &self, 565 Parameters(HotPatchRequest { jump_table }): Parameters<HotPatchRequest>, 566 ) -> Result<CallToolResult, ErrorData> { 567 let bytes = tokio::fs::read(&jump_table).await.map_err(|e| { 568 ErrorData::invalid_params(format!("reading {jump_table}: {e}"), None) 569 })?; 570 let table: subsecond_types::JumpTable = serde_json::from_slice(&bytes).map_err(|e| { 571 ErrorData::invalid_params(format!("parsing {jump_table} as a JumpTable: {e}"), None) 572 })?; 573 let mapped = self 574 .ask(|reply| Request::HotPatch { table, reply }) 575 .await? 576 .map_err(|e| ErrorData::internal_error(e, None))?; 577 Ok(CallToolResult::success(vec![ContentBlock::text(format!( 578 "patched, {mapped} address(es) mapped" 579 ))])) 580 } 581} 582 583impl ServerHandler for Server { 584 fn get_info(&self) -> ServerConfig { 585 ServerConfig::new( 586 ServerCapabilities::builder().enable_tools().enable_tool_list_changed().build(), 587 ) 588 .with_server_info(Implementation::new("jev", "0.0.0")) 589 .with_instructions( 590 "The running jev window: a SNES playing A Link to the Past, played by the zbanks/alttp C bot (the bot tool says what it is doing); \ 591 these tools watch. The dev_mode tool adds the ones that change things. \ 592 Addresses and their meanings are in research/alttp-ram-map.md.", 593 ) 594 } 595 596 async fn list_tools( 597 &self, 598 _: Option<PaginatedRequestParams>, 599 context: RequestContext<RoleServer>, 600 ) -> Result<ListToolsResult, ErrorData> { 601 self.meet(&context.peer); 602 self.heard(|_| ()); 603 let mut tools = self.watching.list_all(); 604 if self.dev_mode.is_on() { 605 tools.extend(self.driving.list_all()); 606 } 607 Ok(ListToolsResult::with_all_items(tools)) 608 } 609 610 fn accepted_subscription_filter(&self, _: &SubscriptionFilter) -> Option<SubscriptionFilter> { 611 Some(SubscriptionFilter::builder().tools_list_changed().build()) 612 } 613 614 async fn listen(&self, context: SubscriptionContext) -> Result<(), ErrorData> { 615 let lock = || self.peers.lock().unwrap_or_else(PoisonError::into_inner); 616 let id = { 617 let mut peers = lock(); 618 let id = peers.next_subscription; 619 peers.next_subscription += 1; 620 peers.all.push(Listener::Subscription { id, sink: context.sink().clone() }); 621 id 622 }; 623 self.heard(|_| ()); 624 context.cancelled().await; 625 lock().all.retain( 626 |listener| !matches!(listener, Listener::Subscription { id: gone, .. } if *gone == id), 627 ); 628 Ok(()) 629 } 630 631 fn get_tool(&self, name: &str) -> Option<Tool> { 632 self.router_of(name)?.get(name).cloned() 633 } 634 635 async fn initialize( 636 &self, 637 request: InitializeRequestParams, 638 context: RequestContext<RoleServer>, 639 ) -> Result<InitializeResult, ErrorData> { 640 context.peer.set_peer_info(request.clone()); 641 self.meet(&context.peer); 642 let client = &request.client_info; 643 self.heard(|activity| { 644 activity.client = Some(format!("{} {}", client.name, client.version)); 645 }); 646 self.negotiate_initialize(&request) 647 } 648 649 /// Every tool call comes through here, so this is where they are logged. 650 async fn call_tool( 651 &self, 652 request: CallToolRequestParams, 653 context: RequestContext<RoleServer>, 654 ) -> Result<CallToolResponse, ErrorData> { 655 self.meet(&context.peer); 656 let tool = request.name.to_string(); 657 let arguments = match &request.arguments { 658 Some(arguments) if !arguments.is_empty() => { 659 serde_json::Value::from(arguments.clone()).to_string() 660 } 661 _ => String::new(), 662 }; 663 let id = self.heard(|activity| { 664 let id = activity.next_id; 665 activity.next_id += 1; 666 if activity.calls.len() == CALLS_KEPT { 667 activity.calls.pop_front(); 668 } 669 activity.calls.push_back(Call { 670 id, 671 at: Instant::now(), 672 tool, 673 arguments, 674 outcome: Outcome::Running, 675 }); 676 id 677 }); 678 679 let response = match self.router_of(&request.name) { 680 Some(router) => router.call(ToolCallContext::new(self, request, context)).await, 681 None => Err(ErrorData::invalid_params( 682 format!("no tool named {:?}", request.name), 683 None, 684 )), 685 }; 686 687 let outcome = match &response { 688 Ok(_) => Outcome::Done, 689 Err(e) => Outcome::Failed(e.message.to_string()), 690 }; 691 self.heard(|activity| { 692 // Gone already if two hundred calls arrived while this one ran. 693 if let Some(call) = activity.calls.iter_mut().find(|call| call.id == id) { 694 call.outcome = outcome; 695 } 696 }); 697 response 698 } 699} 700 701/// Start the server on its own thread and return the two ends the UI thread 702/// reads: what is asked of it, and what there is to show. 703/// 704/// A port that is already taken is reported and the window carries on without 705/// a server: a second window is still a window. 706pub fn serve(ui: egui::Context) -> (mpsc::Receiver<Request>, Arc<Mutex<Activity>>) { 707 let (requests, inbox) = mpsc::channel(); 708 let activity = Arc::new(Mutex::new(Activity::default())); 709 let dev_mode = Arc::clone(&record(&activity).dev_mode); 710 let peers = Peers::default(); 711 let (watching, driving) = (Arc::new(Server::watching()), Arc::new(Server::driving())); 712 713 let spawned = thread::Builder::new().name("mcp".into()).spawn({ 714 let activity = Arc::clone(&activity); 715 move || { 716 let runtime = tokio::runtime::Builder::new_current_thread() 717 .enable_all() 718 .build() 719 .expect("building the tokio runtime"); 720 let served: std::io::Result<()> = runtime.block_on(async { 721 let mut config = StreamableHttpServerConfig::default(); 722 config.session_store = Some(Arc::new(FileSessionStore::open())); 723 let service: StreamableHttpService<Server, LocalSessionManager> = 724 StreamableHttpService::new( 725 { 726 let (ui, activity) = (ui.clone(), Arc::clone(&activity)); 727 move || { 728 Ok(Server { 729 requests: requests.clone(), 730 ui: ui.clone(), 731 activity: Arc::clone(&activity), 732 dev_mode: Arc::clone(&dev_mode), 733 peers: Arc::clone(&peers), 734 known: Arc::default(), 735 watching: Arc::clone(&watching), 736 driving: Arc::clone(&driving), 737 }) 738 } 739 }, 740 Default::default(), 741 config, 742 ); 743 // `/alive` never answers. Whoever asks is holding the line to 744 // know the moment this process is gone, which the MCP client 745 // library hides from them by quietly reconnecting. 746 let alive = axum::routing::get(std::future::pending::<()>); 747 let router = axum::Router::new().route(ALIVE, alive).nest_service("/mcp", service); 748 let listener = tokio::net::TcpListener::bind(ADDRESS).await?; 749 record(&activity).listening = Listening::On; 750 ui.request_repaint(); 751 axum::serve(listener, router).await 752 }); 753 if let Err(e) = served { 754 record(&activity).listening = Listening::Failed(e.to_string()); 755 ui.request_repaint(); 756 } 757 } 758 }); 759 if let Err(e) = spawned { 760 record(&activity).listening = Listening::Failed(e.to_string()); 761 } 762 (inbox, activity) 763}