An MCP server inside the window, so that whoever is developing the app can ask the running console what it is doing instead of photographing it.
The console belongs to the UI thread and stays there. A tool call becomes a
[Request] on a channel, the UI thread answers it at a frame boundary, and
the reply comes back on a oneshot. Nothing here touches the emulator.
Who is developing the app is not who plays the game. The app's own loops
play; a client watches. So the tools come in two sets: the ones that only
look are always there, and the ones that change anything — the pad, the
machine's state, the process — exist only while [DevMode] is on, which it
is not until somebody turns it on. Turning it either way tells every client
that its tool list has changed.
How a client is told depends on which MCP it speaks, and the transport is
streamable HTTP either way. Before 2026-07-28 a client holds a stream open
on its session, so there are sessions, kept in a file ([crate::sessions])
so that a client connected before the app restarted is still connected
after it. From 2026-07-28 there are no sessions and a client asks to be told
with subscriptions/listen, which it makes again itself after a restart.
This crate is the whole native-only side of talking to the app from
outside it: the server, native mcp's stdio proxy ([proxy]) and the
session/snapshot-name files ([states]). It is deliberately its OWN
buck2 target, separate from apps/native's tip crate — see this
directory's CLAUDE.md for why that split exists and what it costs the
hot-patch cycle to get wrong.
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;
Where the server listens. Loopback only: this drives the pad and reads memory, and is for the machine the window is on.
61pub const ADDRESS: &str = "127.0.0.1:7637";
A request here is never answered, only dropped when the window goes.
64pub const ALIVE: &str = "/alive";
The longest a single press may run, in frames: one minute of game time.
A tool call that outlives its client's patience reports nothing at all.
68const LONGEST_PRESS: u32 = 3600;
Work RAM as the SNES addresses it.
71const WRAM_BUS: std::ops::Range<u32> = 0x7E0000..0x800000;
How many tool calls the window keeps to show.
74const CALLS_KEPT: usize = 200;
Whether the server is there to be talked to.
One tool call, as it arrived.
The arguments as the JSON they came in as; empty when there were none.
Whether the tools that change things exist. Off when the window opens, whatever it was when the last one closed.
What the window shows about the server. The server thread writes it and the UI thread reads it.
The client that last introduced itself, by name and version. A client that was already connected before this window started never does.
A panic elsewhere while the record was held leaves a record that is still worth showing, so poisoning is not an error here.
An RGBA picture, rows top to bottom.
What the server asks of the UI thread.
The game's own picture, at the console's resolution.
147 Frame(oneshot::Sender<Picture>),
Offset into work RAM; in range by construction.
The names of the snapshots there are.
157 States(oneshot::Sender<Vec<String>>),
What the zbanks bot says about itself (zbanks::Report), already
serialized to JSON by the app. Read-only.
160 Bot(oneshot::Sender<String>),
Jev's recent goal-choice history (decisions::goal_choice::Chooser::history),
already serialized to JSON by the app, newest first. Read-only.
163 JevHistory(oneshot::Sender<String>),
Hold buttons for frames emulated frames, release them, then reply
with the state the game is left in.
Keep the machine as it stands under a name.
173 SaveState { name: states::Name, reply: oneshot::Sender<Result<(), String>> },
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>> },
Turn the machine off and on again.
177 PowerCycle(oneshot::Sender<alttp::State>),
Keep where the game is, then replace this process with a fresh run of the binary it was started from, which carries on from there.
180 Restart(oneshot::Sender<Result<(), String>>),
Apply a subsecond hot patch. Must run on the UI thread at a frame
boundary, like every other request here - subsecond::apply_patch
rewrites which code the running App value's patch points jump to,
and that jump table is read from the same thread that calls through
it (apps/native/src/app.rs's logic_impl/ui_impl).
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,
subsecond::aslr_reference() in THIS process: main's live address,
which a patch build rebases its jump table against
(research/subsecond-patch-build.md §2/§7).
Somebody to tell that the tool list changed.
271enum Listener {
A session, which lasts until telling it fails.
273 Session(Peer<RoleServer>),
A subscriptions/listen, which lasts until its request is cancelled.
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>>;
One per session.
Set once this session's peer is among peers.
309impl Server {
The router a tool is in, if it exists right now. A driving tool called while dev mode is off is not refused, it is not there.
Any request carries the peer to reach its session by, including the first one from a session restored after a restart, which never introduces itself again.
Note that a client was heard from, record what it said, and have the window redraw to show it.
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}
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 }
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}
Start the server on its own thread and return the two ends the UI thread reads: what is asked of it, and what there is to show.
A port that is already taken is reported and the window carries on without 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}